mirror of https://github.com/apache/cassandra
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
This commit is contained in:
parent
50ee32f685
commit
6faf341044
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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<ReadCommand>
|
||||
|
|
|
|||
|
|
@ -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<IColumn> columns;
|
||||
if (reversed)
|
||||
columns = row.cf.getSortedColumns();
|
||||
else
|
||||
columns = row.cf.getReverseSortedColumns();
|
||||
|
||||
Collection<ByteBuffer> toRemove = new HashSet<ByteBuffer>();
|
||||
|
||||
Iterator<IColumn> 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()
|
||||
{
|
||||
|
|
|
|||
|
|
@ -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<ReadCommand>();
|
||||
commandsToRetry.add(retryCommand);
|
||||
continue;
|
||||
}
|
||||
logger.debug("issuing retry for read command");
|
||||
if (commandsToRetry == Collections.EMPTY_LIST)
|
||||
commandsToRetry = new ArrayList<ReadCommand>();
|
||||
commandsToRetry.add(retryCommand);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (row != null)
|
||||
{
|
||||
command.maybeTrim(row);
|
||||
rows.add(row);
|
||||
}
|
||||
rows.add(row);
|
||||
}
|
||||
}
|
||||
} while (!commandsToRetry.isEmpty());
|
||||
|
|
|
|||
Loading…
Reference in New Issue