diff --git a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java index 80a07cbc38..7c13068ceb 100644 --- a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java +++ b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java @@ -223,15 +223,6 @@ public abstract class AbstractReplicationStrategy return getAddressRanges(temp).get(pendingAddress); } - public QuorumResponseHandler getQuorumResponseHandler(IResponseResolver responseResolver, ConsistencyLevel consistencyLevel) - { - if (consistencyLevel.equals(ConsistencyLevel.LOCAL_QUORUM) || consistencyLevel.equals(ConsistencyLevel.EACH_QUORUM)) - { - return new DatacenterQuorumResponseHandler(responseResolver, consistencyLevel, table); - } - return new QuorumResponseHandler(responseResolver, consistencyLevel, table); - } - public void invalidateCachedTokenEndpointValues() { clearEndpointCache(); diff --git a/src/java/org/apache/cassandra/service/DatacenterQuorumResponseHandler.java b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java similarity index 91% rename from src/java/org/apache/cassandra/service/DatacenterQuorumResponseHandler.java rename to src/java/org/apache/cassandra/service/DatacenterReadCallback.java index 7c6ff9000d..c334d2a6e0 100644 --- a/src/java/org/apache/cassandra/service/DatacenterQuorumResponseHandler.java +++ b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java @@ -37,15 +37,15 @@ import org.apache.cassandra.utils.FBUtilities; /** * Datacenter Quorum response handler blocks for a quorum of responses from the local DC */ -public class DatacenterQuorumResponseHandler extends QuorumResponseHandler +public class DatacenterReadCallback extends ReadCallback { private static final IEndpointSnitch snitch = DatabaseDescriptor.getEndpointSnitch(); private static final String localdc = snitch.getDatacenter(FBUtilities.getLocalAddress()); private AtomicInteger localResponses; - public DatacenterQuorumResponseHandler(IResponseResolver responseResolver, ConsistencyLevel consistencyLevel, String table) + public DatacenterReadCallback(IResponseResolver resolver, ConsistencyLevel consistencyLevel, String table) { - super(responseResolver, consistencyLevel, table); + super(resolver, consistencyLevel, table); localResponses = new AtomicInteger(blockfor); } diff --git a/src/java/org/apache/cassandra/service/QuorumResponseHandler.java b/src/java/org/apache/cassandra/service/ReadCallback.java similarity index 93% rename from src/java/org/apache/cassandra/service/QuorumResponseHandler.java rename to src/java/org/apache/cassandra/service/ReadCallback.java index cf8e75ff69..a6d58f52fd 100644 --- a/src/java/org/apache/cassandra/service/QuorumResponseHandler.java +++ b/src/java/org/apache/cassandra/service/ReadCallback.java @@ -36,9 +36,9 @@ import org.apache.cassandra.thrift.ConsistencyLevel; import org.apache.cassandra.thrift.UnavailableException; import org.apache.cassandra.utils.SimpleCondition; -public class QuorumResponseHandler implements IAsyncCallback +public class ReadCallback implements IAsyncCallback { - protected static final Logger logger = LoggerFactory.getLogger( QuorumResponseHandler.class ); + protected static final Logger logger = LoggerFactory.getLogger( ReadCallback.class ); public final IResponseResolver resolver; protected final SimpleCondition condition = new SimpleCondition(); @@ -48,13 +48,13 @@ public class QuorumResponseHandler implements IAsyncCallback /** * Constructor when response count has to be calculated and blocked for. */ - public QuorumResponseHandler(IResponseResolver resolver, ConsistencyLevel consistencyLevel, String table) + public ReadCallback(IResponseResolver resolver, ConsistencyLevel consistencyLevel, String table) { this.blockfor = determineBlockFor(consistencyLevel, table); this.resolver = resolver; this.startTime = System.currentTimeMillis(); - logger.debug("QuorumResponseHandler blocking for {} responses", blockfor); + logger.debug("ReadCallback blocking for {} responses", blockfor); } public T get() throws TimeoutException, DigestMismatchException, IOException diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 033f5abfcc..e87d2bf811 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -328,7 +328,7 @@ public class StorageProxy implements StorageProxyMBean */ private static List fetchRows(List commands, ConsistencyLevel consistency_level) throws IOException, UnavailableException, TimeoutException { - List> quorumResponseHandlers = new ArrayList>(); + List> readCallbacks = new ArrayList>(); List> commandEndpoints = new ArrayList>(); List rows = new ArrayList(); Set repairs = new HashSet(); @@ -347,7 +347,7 @@ public class StorageProxy implements StorageProxyMBean AbstractReplicationStrategy rs = Table.open(command.table).getReplicationStrategy(); ReadResponseResolver resolver = new ReadResponseResolver(command.table, command.key); - QuorumResponseHandler handler = rs.getQuorumResponseHandler(resolver, consistency_level); + ReadCallback handler = getReadCallback(resolver, command.table, consistency_level); handler.assureSufficientLiveNodes(endpoints); int targets; @@ -374,7 +374,7 @@ public class StorageProxy implements StorageProxyMBean logger.debug("reading " + (m == message ? "data" : "digest") + " for " + command + " from " + m.getMessageId() + "@" + endpoint); } MessagingService.instance().sendRR(messages, endpoints, handler); - quorumResponseHandlers.add(handler); + readCallbacks.add(handler); commandEndpoints.add(endpoints); } @@ -382,22 +382,22 @@ public class StorageProxy implements StorageProxyMBean List> repairResponseHandlers = null; for (int i = 0; i < commands.size(); i++) { - QuorumResponseHandler quorumResponseHandler = quorumResponseHandlers.get(i); + ReadCallback readCallback = readCallbacks.get(i); Row row; ReadCommand command = commands.get(i); List endpoints = commandEndpoints.get(i); try { long startTime2 = System.currentTimeMillis(); - row = quorumResponseHandler.get(); + row = readCallback.get(); if (row != null) rows.add(row); if (logger.isDebugEnabled()) - logger.debug("quorumResponseHandler: " + (System.currentTimeMillis() - startTime2) + " ms."); + logger.debug("Read: " + (System.currentTimeMillis() - startTime2) + " ms."); if (repairs.contains(command)) - repairExecutor.schedule(new RepairRunner(quorumResponseHandler.resolver, command, endpoints), DatabaseDescriptor.getRpcTimeout(), TimeUnit.MILLISECONDS); + repairExecutor.schedule(new RepairRunner(readCallback.resolver, command, endpoints), DatabaseDescriptor.getRpcTimeout(), TimeUnit.MILLISECONDS); } catch (DigestMismatchException ex) { @@ -431,6 +431,15 @@ public class StorageProxy implements StorageProxyMBean return rows; } + static ReadCallback getReadCallback(IResponseResolver resolver, String table, ConsistencyLevel consistencyLevel) + { + if (consistencyLevel.equals(ConsistencyLevel.LOCAL_QUORUM) || consistencyLevel.equals(ConsistencyLevel.EACH_QUORUM)) + { + return new DatacenterReadCallback(resolver, consistencyLevel, table); + } + return new ReadCallback(resolver, consistencyLevel, table); + } + // TODO repair resolver shouldn't take consistencylevel (it should repair exactly as many as it receives replies for) private static RepairCallback repair(ReadCommand command, List endpoints) throws IOException @@ -492,7 +501,7 @@ public class StorageProxy implements StorageProxyMBean // collect replies and resolve according to consistency level RangeSliceResponseResolver resolver = new RangeSliceResponseResolver(command.keyspace, liveEndpoints); AbstractReplicationStrategy rs = Table.open(command.keyspace).getReplicationStrategy(); - QuorumResponseHandler> handler = rs.getQuorumResponseHandler(resolver, consistency_level); + ReadCallback> handler = getReadCallback(resolver, command.keyspace, consistency_level); // TODO bail early if live endpoints can't satisfy requested consistency level for (InetAddress endpoint : liveEndpoints) { @@ -741,7 +750,7 @@ public class StorageProxy implements StorageProxyMBean // collect replies and resolve according to consistency level RangeSliceResponseResolver resolver = new RangeSliceResponseResolver(keyspace, liveEndpoints); AbstractReplicationStrategy rs = Table.open(keyspace).getReplicationStrategy(); - QuorumResponseHandler> handler = rs.getQuorumResponseHandler(resolver, consistency_level); + ReadCallback> handler = getReadCallback(resolver, keyspace, consistency_level); // bail early if live endpoints can't satisfy requested consistency level if(handler.blockfor > liveEndpoints.size()) diff --git a/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java b/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java index 113b59871a..f142231664 100644 --- a/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java +++ b/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java @@ -96,7 +96,7 @@ public class ConsistencyLevelTest extends CleanupHelper IWriteResponseHandler writeHandler = strategy.getWriteResponseHandler(hosts, hintedNodes, c); - QuorumResponseHandler readHandler = strategy.getQuorumResponseHandler(new ReadResponseResolver(table, ByteBufferUtil.bytes("foo")), c); + ReadCallback readHandler = StorageProxy.getReadCallback(new ReadResponseResolver(table, ByteBufferUtil.bytes("foo")), table, c); boolean isWriteUnavailable = false; boolean isReadUnavailable = false;