diff --git a/src/java/org/apache/cassandra/db/ColumnFamily.java b/src/java/org/apache/cassandra/db/ColumnFamily.java index 8b45c697b0..a9efecaac5 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamily.java +++ b/src/java/org/apache/cassandra/db/ColumnFamily.java @@ -397,6 +397,14 @@ public final class ColumnFamily implements IColumnContainer : DatabaseDescriptor.getSubComparator(table, columnFamilyName); } + public static ColumnFamily resolve(ColumnFamily cf1, ColumnFamily cf2) + { + if (cf1 == null) + return cf2; + cf1.resolve(cf2); + return cf1; + } + public void resolve(ColumnFamily cf) { // Row _does_ allow null CF objects :( seems a necessary evil for efficiency diff --git a/src/java/org/apache/cassandra/db/Memtable.java b/src/java/org/apache/cassandra/db/Memtable.java index d61e2845dc..399e37a7c8 100644 --- a/src/java/org/apache/cassandra/db/Memtable.java +++ b/src/java/org/apache/cassandra/db/Memtable.java @@ -157,12 +157,11 @@ public class Memtable implements Comparable, IFlushable { int oldSize = oldCf.size(); int oldObjectCount = oldCf.getColumnCount(); - oldCf.addAll(columnFamily); + oldCf.resolve(columnFamily); int newSize = oldCf.size(); int newObjectCount = oldCf.getColumnCount(); resolveSize(oldSize, newSize); resolveCount(oldObjectCount, newObjectCount); - oldCf.delete(columnFamily); } } diff --git a/src/java/org/apache/cassandra/service/CassandraServer.java b/src/java/org/apache/cassandra/service/CassandraServer.java index 4600b802f2..677efc90bd 100644 --- a/src/java/org/apache/cassandra/service/CassandraServer.java +++ b/src/java/org/apache/cassandra/service/CassandraServer.java @@ -34,6 +34,7 @@ import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.db.filter.QueryPath; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.LogUtil; +import org.apache.cassandra.utils.Pair; import org.apache.thrift.TException; import flexjson.JSONSerializer; @@ -568,24 +569,23 @@ public class CassandraServer implements Cassandra.Iface throw new InvalidRequestException("maxRows must be positive"); } - Map> colMap; // keys are sorted. + List>> rows; try { - colMap = StorageProxy.getRangeSlice(new RangeSliceCommand(keyspace, column_parent, predicate, start_key, finish_key, maxRows)); - if (colMap == null) - throw new RuntimeException("KeySlice list should never be null."); + rows = StorageProxy.getRangeSlice(new RangeSliceCommand(keyspace, column_parent, predicate, start_key, finish_key, maxRows), consistency_level); + assert rows != null; } catch (IOException e) { throw new RuntimeException(e); } - List keySlices = new ArrayList(colMap.size()); - for (String key : colMap.keySet()) + List keySlices = new ArrayList(rows.size()); + for (Pair> row : rows) { - Collection dbList = colMap.get(key); - List svcList = new ArrayList(dbList.size()); - for (org.apache.cassandra.db.IColumn col : dbList) + Collection columns = row.right; + List svcList = new ArrayList(columns.size()); + for (org.apache.cassandra.db.IColumn col : columns) { if (col instanceof org.apache.cassandra.db.Column) svcList.add(new ColumnOrSuperColumn(new org.apache.cassandra.service.Column(col.name(), col.value(), col.timestamp()), null)); @@ -598,7 +598,7 @@ public class CassandraServer implements Cassandra.Iface svcList.add(new ColumnOrSuperColumn(null, new org.apache.cassandra.service.SuperColumn(col.name(), subCols))); } } - keySlices.add(new KeySlice(key, svcList)); + keySlices.add(new KeySlice(row.left, svcList)); } return keySlices; diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index f0f61151ed..6a83f5c061 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -37,6 +37,7 @@ import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.utils.TimedStatsDeque; import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.Pair; import org.apache.cassandra.locator.TokenMetadata; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.gms.FailureDetector; @@ -527,20 +528,17 @@ public class StorageProxy implements StorageProxyMBean return rows; } - static Map> getRangeSlice(RangeSliceCommand rawCommand) throws IOException, UnavailableException, TimedOutException + static List>> getRangeSlice(RangeSliceCommand command, int consistency_level) throws IOException, UnavailableException, TimedOutException { long startTime = System.currentTimeMillis(); TokenMetadata tokenMetadata = StorageService.instance().getTokenMetadata(); - RangeSliceCommand command = rawCommand; InetAddress endPoint = StorageService.instance().findSuitableEndPoint(command.start_key); InetAddress startEndpoint = endPoint; - InetAddress wrapEndpoint = tokenMetadata.getFirstEndpoint(); - TreeSet allRows = new TreeSet(rowComparator); + Map rows = new HashMap(command.max_keys); do { - Message message = command.getMessage(); if (logger.isDebugEnabled()) logger.debug("reading " + command + " from " + message.getMessageId() + "@" + endPoint); @@ -555,44 +553,12 @@ public class StorageProxy implements StorageProxyMBean throw new TimedOutException(); } RangeSliceReply reply = RangeSliceReply.read(responseBody); - List rangeRows = new ArrayList(reply.rows); - - // combine these what what has been seen so far. - if (rangeRows.size() > 0) + for (Row row : reply.rows) { - if (allRows.size() > 0) - { - if (keyComparator.compare(rangeRows.get(rangeRows.size() - 1).key, allRows.first().key) <= 0) - { - // unlikely, but possible - if (rangeRows.get(rangeRows.size() - 1).equals(allRows.first().key)) - { - rangeRows.remove(rangeRows.size() - 1); - } - // put all from rangeRows into allRows. - allRows.addAll(rangeRows); - } - else if (keyComparator.compare(allRows.last().key, rangeRows.get(0).key) <= 0) - { - // common case. deal with simple start/end key overlaps - if (allRows.last().key.equals(rangeRows.get(0))) - { - allRows.remove(allRows.last().key); - } - allRows.addAll(rangeRows); // todo: check logic. - } - else - { - // deal with potential large overlap from scanning the first endpoint, which contains - // both the smallest and largest keys - allRows.addAll(rangeRows); // todo: check logic. - } - } - else - allRows.addAll(rangeRows); // todo: check logic. + rows.put(row.key, ColumnFamily.resolve(row.cf, rows.get(row.key))); } - if (allRows.size() >= rawCommand.max_keys || reply.rangeCompletedLocally) + if (rows.size() >= command.max_keys || reply.rangeCompletedLocally) break; do @@ -600,33 +566,35 @@ public class StorageProxy implements StorageProxyMBean endPoint = tokenMetadata.getSuccessor(endPoint); // TODO move this into the Strategies & modify for RackAwareStrategy } while (!FailureDetector.instance().isAlive(endPoint)); - int maxResults = endPoint == wrapEndpoint ? rawCommand.max_keys : rawCommand.max_keys - allRows.size(); - command = new RangeSliceCommand(command, maxResults); } while (!endPoint.equals(startEndpoint)); - Map> results = new TreeMap>(); - for (Row row : allRows) + List>> results = new ArrayList>>(rows.size()); + for (Map.Entry entry : rows.entrySet()) { - if (row.cf == null) - results.put(row.key, Collections.emptyList()); - else - results.put(row.key, row.cf.getSortedColumns()); + ColumnFamily cf = entry.getValue(); + Collection columns = (cf == null) ? Collections.emptyList() : cf.getSortedColumns(); + results.add(new Pair>(entry.getKey(), columns)); } + Collections.sort(results, new Comparator>>() + { + public int compare(Pair> o1, Pair> o2) + { + return keyComparator.compare(o1.left, o2.left); + } + }); rangeStats.add(System.currentTimeMillis() - startTime); return results; } - static List getKeyRange(RangeCommand rawCommand) throws IOException, UnavailableException, TimedOutException + static List getKeyRange(RangeCommand command) throws IOException, UnavailableException, TimedOutException { long startTime = System.currentTimeMillis(); TokenMetadata tokenMetadata = StorageService.instance().getTokenMetadata(); - List allKeys = new ArrayList(); - RangeCommand command = rawCommand; + Set uniqueKeys = new HashSet(command.maxResults); InetAddress endPoint = StorageService.instance().findSuitableEndPoint(command.startWith); InetAddress startEndpoint = endPoint; - InetAddress wrapEndpoint = tokenMetadata.getFirstEndpoint(); do { @@ -646,49 +614,9 @@ public class StorageProxy implements StorageProxyMBean throw new TimedOutException(); } RangeReply rangeReply = RangeReply.read(responseBody); - List rangeKeys = rangeReply.keys; + uniqueKeys.addAll(rangeReply.keys); - // combine keys from most recent response with the others seen so far - if (rangeKeys.size() > 0) - { - if (allKeys.size() > 0) - { - if (keyComparator.compare(rangeKeys.get(rangeKeys.size() - 1), allKeys.get(0)) <= 0) - { - // unlikely, but possible - if (rangeKeys.get(rangeKeys.size() - 1).equals(allKeys.get(0))) - { - rangeKeys.remove(rangeKeys.size() - 1); - } - rangeKeys.addAll(allKeys); - allKeys = rangeKeys; - } - else if (keyComparator.compare(allKeys.get(allKeys.size() - 1), rangeKeys.get(0)) <= 0) - { - // common case. deal with simple start/end key overlaps - if (allKeys.get(allKeys.size() - 1).equals(rangeKeys.get(0))) - { - allKeys.remove(allKeys.size() - 1); - } - allKeys.addAll(rangeKeys); - } - else - { - // deal with potential large overlap from scanning the first endpoint, which contains - // both the smallest and largest keys - HashSet keys = new HashSet(allKeys); - keys.addAll(rangeKeys); - allKeys = new ArrayList(keys); - Collections.sort(allKeys); - } - } - else - { - allKeys = rangeKeys; - } - } - - if (allKeys.size() >= rawCommand.maxResults || rangeReply.rangeCompletedLocally) + if (uniqueKeys.size() >= command.maxResults || rangeReply.rangeCompletedLocally) { break; } @@ -700,15 +628,15 @@ public class StorageProxy implements StorageProxyMBean // so starting with the largest in our scan of the next node means we'd never see keys from the middle. do { - endPoint = tokenMetadata.getSuccessor(endPoint); // TODO move this into the Strategies & modify for RackAwareStrategy + endPoint = tokenMetadata.getSuccessor(endPoint); } while (!FailureDetector.instance().isAlive(endPoint)); - int maxResults = endPoint.equals(wrapEndpoint) ? rawCommand.maxResults : rawCommand.maxResults - allKeys.size(); - command = new RangeCommand(command.table, command.columnFamily, command.startWith, command.stopAt, maxResults); } while (!endPoint.equals(startEndpoint)); rangeStats.add(System.currentTimeMillis() - startTime); - return (allKeys.size() > rawCommand.maxResults) - ? allKeys.subList(0, rawCommand.maxResults) + List allKeys = new ArrayList(uniqueKeys); + Collections.sort(allKeys, keyComparator); + return (allKeys.size() > command.maxResults) + ? allKeys.subList(0, command.maxResults) : allKeys; } diff --git a/src/java/org/apache/cassandra/utils/Pair.java b/src/java/org/apache/cassandra/utils/Pair.java new file mode 100644 index 0000000000..926ddb3423 --- /dev/null +++ b/src/java/org/apache/cassandra/utils/Pair.java @@ -0,0 +1,34 @@ +package org.apache.cassandra.utils; + +public class Pair +{ + public final T1 left; + public final T2 right; + + public Pair(T1 left, T2 right) + { + this.left = left; + this.right = right; + } + + @Override + public int hashCode() + { + throw new UnsupportedOperationException("todo"); + } + + @Override + public boolean equals(Object obj) + { + throw new UnsupportedOperationException("todo"); + } + + @Override + public String toString() + { + return "Pair(" + + "left=" + left + + ", right=" + right + + ')'; + } +}