From 621899355a23eca6c503aa2ced2e944ff6a5fe66 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Fri, 19 Dec 2014 17:53:46 +0100 Subject: [PATCH] Fix paging with multi-partition IN queries patch by thobbs; reviewed by slebresne for CASSANDRA-8408 --- CHANGES.txt | 1 + .../cql3/statements/SelectStatement.java | 2 +- .../service/pager/MultiPartitionPager.java | 15 ++++----------- .../apache/cassandra/service/pager/Pageable.java | 5 ++++- .../cassandra/service/pager/QueryPagers.java | 7 +------ .../apache/cassandra/service/QueryPagerTest.java | 2 +- 6 files changed, 12 insertions(+), 20 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index bad24e799b..516b4a2da0 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.0.12: + * Fix paging for multi-partition IN queries (CASSANDRA-8408) * Fix MOVED_NODE topology event never being emitted when a node moves its token (CASSANDRA-8373) * Fix validation of indexes in COMPACT tables (CASSANDRA-8156) diff --git a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java index db25716216..f08f6b824a 100644 --- a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java @@ -212,7 +212,7 @@ public class SelectStatement implements CQLStatement, MeasurableForPreparedCache else { List commands = getSliceCommands(variables, limitForQuery, now); - command = commands == null ? null : new Pageable.ReadCommands(commands); + command = commands == null ? null : new Pageable.ReadCommands(commands, limitForQuery); } int pageSize = options.getPageSize(); diff --git a/src/java/org/apache/cassandra/service/pager/MultiPartitionPager.java b/src/java/org/apache/cassandra/service/pager/MultiPartitionPager.java index 35d6752ce9..e478d3a0bb 100644 --- a/src/java/org/apache/cassandra/service/pager/MultiPartitionPager.java +++ b/src/java/org/apache/cassandra/service/pager/MultiPartitionPager.java @@ -46,7 +46,7 @@ class MultiPartitionPager implements QueryPager private int remaining; private int current; - MultiPartitionPager(List commands, ConsistencyLevel consistencyLevel, boolean localQuery, PagingState state) + MultiPartitionPager(List commands, ConsistencyLevel consistencyLevel, boolean localQuery, PagingState state, int limitForQuery) { int i = 0; // If it's not the beginning (state != null), we need to find where we were and skip previous commands @@ -76,7 +76,8 @@ class MultiPartitionPager implements QueryPager throw new IllegalArgumentException("All commands must have the same timestamp or weird results may happen."); pagers[j - i] = makePager(command, consistencyLevel, localQuery, null); } - remaining = state == null ? computeRemaining(pagers) : state.remaining; + + remaining = state == null ? limitForQuery : state.remaining; } private static SinglePartitionPager makePager(ReadCommand command, ConsistencyLevel consistencyLevel, boolean localQuery, PagingState state) @@ -86,14 +87,6 @@ class MultiPartitionPager implements QueryPager : new NamesQueryPager((SliceByNamesReadCommand)command, consistencyLevel, localQuery); } - private static int computeRemaining(SinglePartitionPager[] pagers) - { - long remaining = 0; - for (SinglePartitionPager pager : pagers) - remaining += pager.maxRemaining(); - return remaining > Integer.MAX_VALUE ? Integer.MAX_VALUE : (int)remaining; - } - public PagingState state() { // Sets current to the first non-exhausted pager @@ -123,7 +116,7 @@ class MultiPartitionPager implements QueryPager { List result = new ArrayList(); - int remainingThisQuery = pageSize; + int remainingThisQuery = Math.min(remaining, pageSize); while (remainingThisQuery > 0 && !isExhausted()) { // isExhausted has set us on the first non-exhausted pager diff --git a/src/java/org/apache/cassandra/service/pager/Pageable.java b/src/java/org/apache/cassandra/service/pager/Pageable.java index 3a69bf4677..d4986f7acc 100644 --- a/src/java/org/apache/cassandra/service/pager/Pageable.java +++ b/src/java/org/apache/cassandra/service/pager/Pageable.java @@ -30,9 +30,12 @@ public interface Pageable { public final List commands; - public ReadCommands(List commands) + public final int limitForQuery; + + public ReadCommands(List commands, int limitForQuery) { this.commands = commands; + this.limitForQuery = limitForQuery; } } } diff --git a/src/java/org/apache/cassandra/service/pager/QueryPagers.java b/src/java/org/apache/cassandra/service/pager/QueryPagers.java index 65112aaf77..72d76fe84f 100644 --- a/src/java/org/apache/cassandra/service/pager/QueryPagers.java +++ b/src/java/org/apache/cassandra/service/pager/QueryPagers.java @@ -98,7 +98,7 @@ public class QueryPagers if (commands.size() == 1) return pager(commands.get(0), consistencyLevel, local, state); - return new MultiPartitionPager(commands, consistencyLevel, local, state); + return new MultiPartitionPager(commands, consistencyLevel, local, state, ((Pageable.ReadCommands) command).limitForQuery); } else if (command instanceof ReadCommand) { @@ -115,11 +115,6 @@ public class QueryPagers } } - public static QueryPager pager(Pageable command, ConsistencyLevel consistencyLevel) - { - return pager(command, consistencyLevel, false, null); - } - public static QueryPager pager(Pageable command, ConsistencyLevel consistencyLevel, PagingState state) { return pager(command, consistencyLevel, false, state); diff --git a/test/unit/org/apache/cassandra/service/QueryPagerTest.java b/test/unit/org/apache/cassandra/service/QueryPagerTest.java index 0645433de2..7dbd7b97c6 100644 --- a/test/unit/org/apache/cassandra/service/QueryPagerTest.java +++ b/test/unit/org/apache/cassandra/service/QueryPagerTest.java @@ -236,7 +236,7 @@ public class QueryPagerTest extends SchemaLoader QueryPager pager = QueryPagers.localPager(new Pageable.ReadCommands(new ArrayList() {{ add(sliceQuery("k1", "c2", "c6", 10)); add(sliceQuery("k4", "c3", "c5", 10)); - }})); + }}, 10)); List page;