From 93f71af6b3fedd7a4f35804df6a150aaae66ed64 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Tue, 27 Jul 2010 03:54:22 +0000 Subject: [PATCH] merge from 0.6 git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@979514 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 +- .../service/ReadResponseResolver.java | 4 +- .../cassandra/service/StorageProxy.java | 46 ++++++++++++------- 3 files changed, 32 insertions(+), 20 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 1d416bdcc7..722385507d 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -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 diff --git a/src/java/org/apache/cassandra/service/ReadResponseResolver.java b/src/java/org/apache/cassandra/service/ReadResponseResolver.java index d4dcc04dce..50bd7d2207 100644 --- a/src/java/org/apache/cassandra/service/ReadResponseResolver.java +++ b/src/java/org/apache/cassandra/service/ReadResponseResolver.java @@ -60,8 +60,8 @@ public class ReadResponseResolver implements IResponseResolver public Row resolve(Collection responses) throws DigestMismatchException, IOException { long startTime = System.currentTimeMillis(); - List versions = new ArrayList(); - List endpoints = new ArrayList(); + List versions = new ArrayList(responses.size()); + List endpoints = new ArrayList(responses.size()); DecoratedKey key = null; byte[] digest = new byte[0]; boolean isDigestQuery = false; diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 582829d869..5a40d8db1f 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -395,8 +395,7 @@ public class StorageProxy implements StorageProxyMBean List commandEndpoints = new ArrayList(); List rows = new ArrayList(); - 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 quorumResponseHandler: quorumResponseHandlers) + // read results and make a second pass for any digest mismatches + List> repairResponseHandlers = null; + for (int i = 0; i < commands.size(); i++) { + QuorumResponseHandler 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 quorumResponseHandlerRepair = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), ConsistencyLevel.QUORUM, command.table); + QuorumResponseHandler 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>(); + repairResponseHandlers.add(qrhRepair); + } + } + } + + // read the results for the digest mismatch retries + if (repairResponseHandlers != null) + { + for (QuorumResponseHandler 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;