From 6faf3410447ae81f5b77ce0570d6dbcb99bc8abf Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Fri, 4 Nov 2011 08:35:33 +0000 Subject: [PATCH] Never return more columns than requested patch by byronclark; reviewed by slebresne for CASSANDRA-3303 and 3395 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0@1197426 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../org/apache/cassandra/db/ReadCommand.java | 17 ++++- .../cassandra/db/SliceFromReadCommand.java | 66 +++++++++++++++++++ .../cassandra/service/StorageProxy.java | 43 +++++------- 4 files changed, 98 insertions(+), 29 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index b5a3e3594d..9403627040 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -11,6 +11,7 @@ * add JMX call to clean (failed) repair sessions (CASSANDRA-3316) * fix sstableloader reference acquisition bug (CASSANDRA-3438) * fix estimated row size regression (CASSANDRA-3451) + * make sure we don't return more columns than asked (CASSANDRA-3303, 3395) Merged from 0.8: * acquire compactionlock during truncate (CASSANDRA-3399) * fix displaying cfdef entries for super columnfamilies (CASSANDRA-3415) diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index 7a1dc6a30f..5ebb244f07 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -31,6 +31,7 @@ import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessageProducer; import org.apache.cassandra.service.IReadCommand; +import org.apache.cassandra.service.RepairCallback; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.FBUtilities; @@ -66,7 +67,7 @@ public abstract class ReadCommand implements MessageProducer, IReadCommand this.queryPath = queryPath; this.commandType = cmdType; } - + public boolean isDigestQuery() { return isDigestQuery; @@ -81,7 +82,7 @@ public abstract class ReadCommand implements MessageProducer, IReadCommand { return queryPath.columnFamilyName; } - + public abstract ReadCommand copy(); public abstract Row getRow(Table table) throws IOException; @@ -95,6 +96,18 @@ public abstract class ReadCommand implements MessageProducer, IReadCommand { return table; } + + // maybeGenerateRetryCommand is used to generate a retry for short reads + public ReadCommand maybeGenerateRetryCommand(RepairCallback handler, Row row) + { + return null; + } + + // maybeTrim removes columns from a response that is too long + public void maybeTrim(Row row) + { + // noop + } } class ReadCommandSerializer implements IVersionedSerializer diff --git a/src/java/org/apache/cassandra/db/SliceFromReadCommand.java b/src/java/org/apache/cassandra/db/SliceFromReadCommand.java index d9d726366e..51c1602efa 100644 --- a/src/java/org/apache/cassandra/db/SliceFromReadCommand.java +++ b/src/java/org/apache/cassandra/db/SliceFromReadCommand.java @@ -19,17 +19,25 @@ package org.apache.cassandra.db; import java.io.*; import java.nio.ByteBuffer; +import java.util.Collection; +import java.util.HashSet; +import java.util.Iterator; import org.apache.cassandra.db.filter.QueryFilter; import org.apache.cassandra.db.filter.QueryPath; import org.apache.cassandra.io.IVersionedSerializer; +import org.apache.cassandra.service.RepairCallback; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.thrift.ColumnParent; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; public class SliceFromReadCommand extends ReadCommand { + static final Logger logger = LoggerFactory.getLogger(SliceFromReadCommand.class); + public final ByteBuffer start, finish; public final boolean reversed; public final int count; @@ -61,6 +69,64 @@ public class SliceFromReadCommand extends ReadCommand return table.getRow(QueryFilter.getSliceFilter(dk, queryPath, start, finish, reversed, count)); } + @Override + public ReadCommand maybeGenerateRetryCommand(RepairCallback handler, Row row) + { + int maxLiveColumns = handler.getMaxLiveColumns(); + int liveColumnsInRow = row != null ? row.cf.getLiveColumnCount() : 0; + + assert maxLiveColumns <= count; + if ((maxLiveColumns == count) && (liveColumnsInRow < count)) + { + int retryCount = count + count - liveColumnsInRow; + return new RetriedSliceFromReadCommand(table, key, queryPath, start, finish, reversed, count, retryCount); + } + + return null; + } + + @Override + public void maybeTrim(Row row) + { + if ((row == null) || (row.cf == null)) + return; + + int liveColumnsInRow = row.cf.getLiveColumnCount(); + + if (liveColumnsInRow > getRequestedCount()) + { + int columnsToTrim = liveColumnsInRow - getRequestedCount(); + + logger.debug("trimming {} live columns to the originally requested {}", row.cf.getLiveColumnCount(), getRequestedCount()); + + Collection columns; + if (reversed) + columns = row.cf.getSortedColumns(); + else + columns = row.cf.getReverseSortedColumns(); + + Collection toRemove = new HashSet(); + + Iterator columnIterator = columns.iterator(); + while (columnIterator.hasNext() && (toRemove.size() < columnsToTrim)) + { + IColumn column = columnIterator.next(); + if (column.isLive()) + toRemove.add(column.name()); + } + + for (ByteBuffer columnName : toRemove) + { + row.cf.remove(columnName); + } + } + } + + protected int getRequestedCount() + { + return count; + } + @Override public String toString() { diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 1e31348425..1c056cf3f0 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -685,7 +685,10 @@ public class StorageProxy implements StorageProxyMBean long startTime2 = System.currentTimeMillis(); Row row = handler.get(); if (row != null) + { + command.maybeTrim(row); rows.add(row); + } if (logger.isDebugEnabled()) logger.debug("Read: " + (System.currentTimeMillis() - startTime2) + " ms."); @@ -739,35 +742,21 @@ public class StorageProxy implements StorageProxyMBean throw new AssertionError(e); // full data requested from each node here, no digests should be sent } - // retry short reads, otherwise add the row to our resultset - if (command instanceof SliceFromReadCommand) + ReadCommand retryCommand = command.maybeGenerateRetryCommand(handler, row); + if (retryCommand != null) { - // short reads are only possible on SliceFromReadCommand - SliceFromReadCommand sliceCommand = (SliceFromReadCommand) command; - int maxLiveColumns = handler.getMaxLiveColumns(); - int liveColumnsInRow = row != null ? row.cf.getLiveColumnCount() : 0; - - assert maxLiveColumns <= sliceCommand.count; - if ((maxLiveColumns == sliceCommand.count) && (liveColumnsInRow < sliceCommand.count)) - { - logger.debug("detected short read: expected {} columns, but only resolved {} columns", - sliceCommand.count, liveColumnsInRow); - - int retryCount = sliceCommand.count + sliceCommand.count - liveColumnsInRow; - SliceFromReadCommand retryCommand = new SliceFromReadCommand(command.table, - command.key, - command.queryPath, - sliceCommand.start, - sliceCommand.finish, - sliceCommand.reversed, - retryCount); - if (commandsToRetry == Collections.EMPTY_LIST) - commandsToRetry = new ArrayList(); - commandsToRetry.add(retryCommand); - continue; - } + logger.debug("issuing retry for read command"); + if (commandsToRetry == Collections.EMPTY_LIST) + commandsToRetry = new ArrayList(); + commandsToRetry.add(retryCommand); + continue; + } + + if (row != null) + { + command.maybeTrim(row); + rows.add(row); } - rows.add(row); } } } while (!commandsToRetry.isEmpty());