merge from 0.6

git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@979514 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
Jonathan Ellis 2010-07-27 03:54:22 +00:00
parent 034845d8a3
commit 93f71af6b3
3 changed files with 32 additions and 20 deletions

View File

@ -69,7 +69,7 @@ dev
when determining whether to do local read for CL.ONE (CASSANDRA-1317)
* fix read repair to use requested consistency level on digest mismatch,
rather than assuming QUORUM (CASSANDRA-1316)
* process digest mismatch re-reads in parallel (CASSANDRA-1323)
0.6.3

View File

@ -60,8 +60,8 @@ public class ReadResponseResolver implements IResponseResolver<Row>
public Row resolve(Collection<Message> responses) throws DigestMismatchException, IOException
{
long startTime = System.currentTimeMillis();
List<ColumnFamily> versions = new ArrayList<ColumnFamily>();
List<InetAddress> endpoints = new ArrayList<InetAddress>();
List<ColumnFamily> versions = new ArrayList<ColumnFamily>(responses.size());
List<InetAddress> endpoints = new ArrayList<InetAddress>(responses.size());
DecoratedKey key = null;
byte[] digest = new byte[0];
boolean isDigestQuery = false;

View File

@ -395,8 +395,7 @@ public class StorageProxy implements StorageProxyMBean
List<InetAddress[]> commandEndpoints = new ArrayList<InetAddress[]>();
List<Row> rows = new ArrayList<Row>();
int commandIndex = 0;
// send out read requests
for (ReadCommand command: commands)
{
assert !command.isDigestQuery();
@ -428,10 +427,13 @@ public class StorageProxy implements StorageProxyMBean
commandEndpoints.add(endpoints);
}
for (QuorumResponseHandler<Row> quorumResponseHandler: quorumResponseHandlers)
// read results and make a second pass for any digest mismatches
List<QuorumResponseHandler<Row>> repairResponseHandlers = null;
for (int i = 0; i < commands.size(); i++)
{
QuorumResponseHandler<Row> quorumResponseHandler = quorumResponseHandlers.get(i);
Row row;
ReadCommand command = commands.get(commandIndex);
ReadCommand command = commands.get(i);
try
{
long startTime2 = System.currentTimeMillis();
@ -447,24 +449,34 @@ public class StorageProxy implements StorageProxyMBean
if (randomlyReadRepair(command))
{
AbstractReplicationStrategy rs = StorageService.instance.getReplicationStrategy(command.table);
QuorumResponseHandler<Row> quorumResponseHandlerRepair = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), ConsistencyLevel.QUORUM, command.table);
QuorumResponseHandler<Row> qrhRepair = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), ConsistencyLevel.QUORUM, command.table);
if (logger.isDebugEnabled())
logger.debug("Digest mismatch:", ex);
Message messageRepair = command.makeReadMessage();
MessagingService.instance.sendRR(messageRepair, commandEndpoints.get(commandIndex), quorumResponseHandlerRepair);
try
{
row = quorumResponseHandlerRepair.get();
if (row != null)
rows.add(row);
}
catch (DigestMismatchException e)
{
throw new AssertionError(e); // full data requested from each node here, no digests should be sent
}
MessagingService.instance.sendRR(messageRepair, commandEndpoints.get(i), qrhRepair);
if (repairResponseHandlers == null)
repairResponseHandlers = new ArrayList<QuorumResponseHandler<Row>>();
repairResponseHandlers.add(qrhRepair);
}
}
}
// read the results for the digest mismatch retries
if (repairResponseHandlers != null)
{
for (QuorumResponseHandler<Row> handler : repairResponseHandlers)
{
try
{
Row row = handler.get();
if (row != null)
rows.add(row);
}
catch (DigestMismatchException e)
{
throw new AssertionError(e); // full data requested from each node here, no digests should be sent
}
}
commandIndex++;
}
return rows;