From cbee7a7fe00083595754a4e48271305cc9c7b42c Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 25 Jun 2010 22:11:54 +0000 Subject: [PATCH] replace sorting of unwrapped range in token order, fixing a regression introduced in r948934. patch by jbellis; reviewed by eevans for CASSANDRA-1198 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.6@958133 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 - .../cassandra/service/StorageProxy.java | 62 +++++++++++-------- 2 files changed, 36 insertions(+), 27 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 07185e938f..0cc3dbba57 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,6 +1,5 @@ 0.6.3 * retry to make streaming connections up to 8 times. (CASSANDRA-1019) - * fix potential for duplicate rows seen by Hadoop jobs (CASSANDRA-1042) * reject describe_ring() calls on invalid keyspaces (CASSANDRA-1111) * fix cache size calculation for size of 100% (CASSANDRA-1129) * fix cache capacity only being recalculated once (CASSANDRA-1129) diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index fc29f0c092..3c7178f2a7 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -557,14 +557,18 @@ public class StorageProxy implements StorageProxyMBean final String table = command.keyspace; int responseCount = determineBlockFor(DatabaseDescriptor.getReplicationFactor(table), consistency_level); - List>> ranges = getRestrictedRanges(command.range, command.keyspace, responseCount); + List ranges = getRestrictedRanges(command.range, command.keyspace, responseCount); // now scan until we have enough results List rows = new ArrayList(command.max_keys); - for (Pair> pair : getRangeIterator(ranges, command.range.left)) + for (AbstractBounds range : getRangeIterator(ranges, command.range.left)) { - AbstractBounds range = pair.left; - List endpoints = pair.right; + List liveEndpoints = StorageService.instance.getLiveNaturalEndpoints(command.keyspace, range.right); + if (liveEndpoints.size() < responseCount) + throw new UnavailableException(); + DatabaseDescriptor.getEndPointSnitch(command.keyspace).sortByProximity(FBUtilities.getLocalAddress(), liveEndpoints); + List endpoints = liveEndpoints.subList(0, responseCount); + RangeSliceCommand c2 = new RangeSliceCommand(command.keyspace, command.column_family, command.super_column, command.predicate, range, command.max_keys); Message message = c2.getMessage(); @@ -607,30 +611,30 @@ public class StorageProxy implements StorageProxyMBean /** * returns an iterator that will return ranges in ring order, starting with the one that contains the start token */ - private static Iterable>> getRangeIterator(final List>> ranges, Token start) + private static Iterable getRangeIterator(final List ranges, Token start) { // find the one to start with int i; for (i = 0; i < ranges.size(); i++) { - AbstractBounds range = ranges.get(i).left; + AbstractBounds range = ranges.get(i); if (range.contains(start) || range.left.equals(start)) break; } - AbstractBounds range = ranges.get(i).left; + AbstractBounds range = ranges.get(i); assert range.contains(start) || range.left.equals(start); // make sure the loop didn't just end b/c ranges were exhausted // return an iterable that starts w/ the correct range and iterates the rest in ring order final int begin = i; - return new Iterable>>() + return new Iterable() { - public Iterator>> iterator() + public Iterator iterator() { - return new AbstractIterator>>() + return new AbstractIterator() { int n = 0; - protected Pair> computeNext() + protected AbstractBounds computeNext() { if (n == ranges.size()) return endOfData(); @@ -655,33 +659,39 @@ public class StorageProxy implements StorageProxyMBean * D, but we don't want any other results from it until after the (D, T] range. Unwrapping so that * the ranges we consider are (D, T], (T, MIN], (MIN, D] fixes this. */ - private static List>> getRestrictedRanges(AbstractBounds queryRange, String keyspace, int responseCount) + private static List getRestrictedRanges(AbstractBounds queryRange, String keyspace, int responseCount) throws UnavailableException { TokenMetadata tokenMetadata = StorageService.instance.getTokenMetadata(); - Iterator iter = TokenMetadata.ringIterator(tokenMetadata.sortedTokens(), queryRange.left); - List>> ranges = new ArrayList>>(); - while (iter.hasNext()) - { - Token nodeToken = iter.next(); - Range nodeRange = new Range(tokenMetadata.getPredecessor(nodeToken), nodeToken); - List endpoints = StorageService.instance.getLiveNaturalEndpoints(keyspace, nodeToken); - if (endpoints.size() < responseCount) - throw new UnavailableException(); - DatabaseDescriptor.getEndPointSnitch(keyspace).sortByProximity(FBUtilities.getLocalAddress(), endpoints); - List endpointsForCL = endpoints.subList(0, responseCount); - Set restrictedRanges = queryRange.restrictTo(nodeRange); - for (AbstractBounds range : restrictedRanges) + List ranges = new ArrayList(); + // for each node, compute its intersection with the query range, and add its unwrapped components to our list + for (Token nodeToken : tokenMetadata.sortedTokens()) + { + Range nodeRange = new Range(tokenMetadata.getPredecessor(nodeToken), nodeToken); + for (AbstractBounds range : queryRange.restrictTo(nodeRange)) { for (AbstractBounds unwrapped : range.unwrap()) { if (logger.isDebugEnabled()) logger.debug("Adding to restricted ranges " + unwrapped + " for " + nodeRange); - ranges.add(new Pair>(unwrapped, endpointsForCL)); + ranges.add(unwrapped); } } } + + // re-sort ranges in ring order, post-unwrapping + Comparator comparator = new Comparator() + { + public int compare(AbstractBounds o1, AbstractBounds o2) + { + // no restricted ranges will overlap so we don't need to worry about inclusive vs exclusive left, + // just sort by raw token position. + return o1.left.compareTo(o2.left); + } + }; + Collections.sort(ranges, comparator); + return ranges; }