diff --git a/src/java/org/apache/cassandra/db/filter/DataLimits.java b/src/java/org/apache/cassandra/db/filter/DataLimits.java index f2d250879f41..5b56eda1e13e 100644 --- a/src/java/org/apache/cassandra/db/filter/DataLimits.java +++ b/src/java/org/apache/cassandra/db/filter/DataLimits.java @@ -720,7 +720,7 @@ public boolean isUnlimited() public DataLimits forShortReadRetry(int toFetch) { - return new CQLLimits(toFetch); + return new CQLGroupByLimits(toFetch, groupPerPartitionLimit, rowLimit, groupBySpec, state); } @Override diff --git a/src/java/org/apache/cassandra/service/reads/ShortReadPartitionsProtection.java b/src/java/org/apache/cassandra/service/reads/ShortReadPartitionsProtection.java index d1562b1b4de9..5b310a813205 100644 --- a/src/java/org/apache/cassandra/service/reads/ShortReadPartitionsProtection.java +++ b/src/java/org/apache/cassandra/service/reads/ShortReadPartitionsProtection.java @@ -150,7 +150,7 @@ public UnfilteredPartitionIterator moreContents() * then future ShortReadRowsProtection.moreContents() calls will fetch the missing ones. */ int toQuery = command.limits().count() != DataLimits.NO_LIMIT - ? command.limits().count() - mergedResultCounter.rowsCounted() + ? Math.max(0, command.limits().count() - mergedResultCounter.counted()) : command.limits().perPartitionCount(); ColumnFamilyStore.metricsFor(command.metadata().id).shortReadProtectionRequests.mark();