From 6d96f53add0c0fc9193a0dfac8c42ac18a6a97e7 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Sat, 5 Dec 2009 00:24:19 +0000 Subject: [PATCH] move "are we done yet" check in quorum read entirely into RRR.isDataPresent instead of partly there and partly in QRH.response patch by jbellis; reviewed by Stu Hood for CASSANDRA-568 git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@887465 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/service/ConsistencyManager.java | 5 ++--- .../service/QuorumResponseHandler.java | 7 +------ .../service/ReadResponseResolver.java | 18 ++++++++++-------- .../apache/cassandra/service/StorageProxy.java | 4 ++-- 4 files changed, 15 insertions(+), 19 deletions(-) diff --git a/src/java/org/apache/cassandra/service/ConsistencyManager.java b/src/java/org/apache/cassandra/service/ConsistencyManager.java index cb96f157d1..39a4e5200c 100644 --- a/src/java/org/apache/cassandra/service/ConsistencyManager.java +++ b/src/java/org/apache/cassandra/service/ConsistencyManager.java @@ -79,10 +79,9 @@ class ConsistencyManager implements Runnable private void doReadRepair() throws IOException { - IResponseResolver readResponseResolver = new ReadResponseResolver(table_); - /* Add the local storage endpoint to the replicas_ list */ replicas_.add(FBUtilities.getLocalAddress()); - IAsyncCallback responseHandler = new DataRepairHandler(ConsistencyManager.this.replicas_.size(), readResponseResolver); + IResponseResolver readResponseResolver = new ReadResponseResolver(table_, replicas_.size()); + IAsyncCallback responseHandler = new DataRepairHandler(replicas_.size(), readResponseResolver); ReadCommand readCommand = constructReadMessage(false); Message message = readCommand.makeReadMessage(); if (logger_.isDebugEnabled()) diff --git a/src/java/org/apache/cassandra/service/QuorumResponseHandler.java b/src/java/org/apache/cassandra/service/QuorumResponseHandler.java index 7fcdbbab1d..7b0975dbe3 100644 --- a/src/java/org/apache/cassandra/service/QuorumResponseHandler.java +++ b/src/java/org/apache/cassandra/service/QuorumResponseHandler.java @@ -36,17 +36,12 @@ public class QuorumResponseHandler implements IAsyncCallback { protected static final Logger logger = Logger.getLogger( QuorumResponseHandler.class ); protected final SimpleCondition condition = new SimpleCondition(); - private final int responseCount; protected final List responses; private IResponseResolver responseResolver; private final long startTime; public QuorumResponseHandler(int responseCount, IResponseResolver responseResolver) { - assert 1 <= responseCount && responseCount <= DatabaseDescriptor.getReplicationFactor() - : "invalid response count " + responseCount; - - this.responseCount = responseCount; responses = new ArrayList(responseCount); this.responseResolver = responseResolver; startTime = System.currentTimeMillis(); @@ -94,7 +89,7 @@ public class QuorumResponseHandler implements IAsyncCallback return; responses.add(message); - if (responses.size() >= responseCount && responseResolver.isDataPresent(responses)) + if (responseResolver.isDataPresent(responses)) { condition.signal(); } diff --git a/src/java/org/apache/cassandra/service/ReadResponseResolver.java b/src/java/org/apache/cassandra/service/ReadResponseResolver.java index 9863d74277..ab3c537597 100644 --- a/src/java/org/apache/cassandra/service/ReadResponseResolver.java +++ b/src/java/org/apache/cassandra/service/ReadResponseResolver.java @@ -33,23 +33,22 @@ import java.net.InetAddress; import org.apache.cassandra.net.Message; import org.apache.cassandra.utils.LogUtil; import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.log4j.Logger; - -/** - * This class is used by all read functions and is called by the Quorum - * when at least a few of the servers (few is specified in Quorum) - * have sent the response . The resolve function then schedules read repair - * and resolution of read data from the various servers. - */ public class ReadResponseResolver implements IResponseResolver { private static Logger logger_ = Logger.getLogger(ReadResponseResolver.class); private final String table; + private final int responseCount; - public ReadResponseResolver(String table) + public ReadResponseResolver(String table, int responseCount) { + assert 1 <= responseCount && responseCount <= DatabaseDescriptor.getReplicationFactor() + : "invalid response count " + responseCount; + + this.responseCount = responseCount; this.table = table; } @@ -152,6 +151,9 @@ public class ReadResponseResolver implements IResponseResolver public boolean isDataPresent(List responses) { + if (responses.size() < responseCount) + return false; + boolean isDataPresent = false; for (Message response : responses) { diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 6a83f5c061..0ce764fcd3 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -438,7 +438,7 @@ public class StorageProxy implements StorageProxyMBean } if (n < DatabaseDescriptor.getQuorum()) throw new UnavailableException(); - QuorumResponseHandler quorumResponseHandler = new QuorumResponseHandler(DatabaseDescriptor.getQuorum(), new ReadResponseResolver(command.table)); + QuorumResponseHandler quorumResponseHandler = new QuorumResponseHandler(DatabaseDescriptor.getQuorum(), new ReadResponseResolver(command.table, DatabaseDescriptor.getQuorum())); MessagingService.instance().sendRR(messages, endPoints, quorumResponseHandler); quorumResponseHandlers.add(quorumResponseHandler); commandEndPoints.add(endPoints); @@ -466,7 +466,7 @@ public class StorageProxy implements StorageProxyMBean { if (DatabaseDescriptor.getConsistencyCheck()) { - IResponseResolver readResponseResolverRepair = new ReadResponseResolver(command.table); + IResponseResolver readResponseResolverRepair = new ReadResponseResolver(command.table, DatabaseDescriptor.getQuorum()); QuorumResponseHandler quorumResponseHandlerRepair = new QuorumResponseHandler( DatabaseDescriptor.getQuorum(), readResponseResolverRepair);