From fb07f2fae9d30e95ee599b4b993dfe34dfe408a2 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Tue, 15 Feb 2011 20:50:21 +0000 Subject: [PATCH] revert #2069 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7.2@1071046 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 - .../apache/cassandra/concurrent/Stage.java | 4 +- .../cassandra/concurrent/StageManager.java | 1 - .../cassandra/db/RangeSliceCommand.java | 8 +- .../org/apache/cassandra/db/ReadCommand.java | 8 +- .../apache/cassandra/dht/BootStrapper.java | 5 - .../org/apache/cassandra/net/AsyncResult.java | 5 - .../cassandra/net/IMessageCallback.java | 5 - .../cassandra/net/MessagingService.java | 2 +- .../service/AbstractRowResolver.java | 70 ----- .../service/AsyncRepairCallback.java | 41 --- .../service/DatacenterReadCallback.java | 17 +- .../DatacenterSyncWriteResponseHandler.java | 5 - .../cassandra/service/IReadCommand.java | 6 - .../service/RangeSliceResponseResolver.java | 5 +- .../cassandra/service/ReadCallback.java | 105 +------ .../service/ReadResponseResolver.java | 257 ++++++++++++++++++ .../cassandra/service/RepairCallback.java | 17 +- .../cassandra/service/RowDigestResolver.java | 125 --------- .../cassandra/service/RowRepairResolver.java | 148 ---------- .../cassandra/service/StorageProxy.java | 156 ++++++++--- .../service/TruncateResponseHandler.java | 5 - .../service/WriteResponseHandler.java | 5 - .../service/ConsistencyLevelTest.java | 14 +- ...est.java => ReadResponseResolverTest.java} | 12 +- 25 files changed, 404 insertions(+), 623 deletions(-) delete mode 100644 src/java/org/apache/cassandra/service/AbstractRowResolver.java delete mode 100644 src/java/org/apache/cassandra/service/AsyncRepairCallback.java delete mode 100644 src/java/org/apache/cassandra/service/IReadCommand.java delete mode 100644 src/java/org/apache/cassandra/service/RowDigestResolver.java delete mode 100644 src/java/org/apache/cassandra/service/RowRepairResolver.java rename test/unit/org/apache/cassandra/service/{RowResolverTest.java => ReadResponseResolverTest.java} (84%) diff --git a/CHANGES.txt b/CHANGES.txt index faa7dc6845..8ca896f044 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -2,7 +2,6 @@ * copy DecoratedKey.key when inserting into caches to avoid retaining a reference to the underlying buffer (CASSANDRA-2102) * format subcolumn names with subcomparator (CASSANDRA-2136) - * lower-latency read repair (CASSANDRA-2069) * fix column bloom filter deserialization (CASSANDRA-2165) diff --git a/src/java/org/apache/cassandra/concurrent/Stage.java b/src/java/org/apache/cassandra/concurrent/Stage.java index 4bd2f367f8..924f413ac5 100644 --- a/src/java/org/apache/cassandra/concurrent/Stage.java +++ b/src/java/org/apache/cassandra/concurrent/Stage.java @@ -31,8 +31,7 @@ public enum Stage ANTI_ENTROPY, MIGRATION, MISC, - INTERNAL_RESPONSE, - READ_REPAIR; + INTERNAL_RESPONSE; public String getJmxType() { @@ -48,7 +47,6 @@ public enum Stage case MUTATION: case READ: case REQUEST_RESPONSE: - case READ_REPAIR: return "request"; default: throw new AssertionError("Unknown stage " + this); diff --git a/src/java/org/apache/cassandra/concurrent/StageManager.java b/src/java/org/apache/cassandra/concurrent/StageManager.java index 6cf44e9512..e4a0a7d3e5 100644 --- a/src/java/org/apache/cassandra/concurrent/StageManager.java +++ b/src/java/org/apache/cassandra/concurrent/StageManager.java @@ -50,7 +50,6 @@ public class StageManager stages.put(Stage.ANTI_ENTROPY, new JMXEnabledThreadPoolExecutor(Stage.ANTI_ENTROPY)); stages.put(Stage.MIGRATION, new JMXEnabledThreadPoolExecutor(Stage.MIGRATION)); stages.put(Stage.MISC, new JMXEnabledThreadPoolExecutor(Stage.MISC)); - stages.put(Stage.READ_REPAIR, multiThreadedStage(Stage.READ_REPAIR, Runtime.getRuntime().availableProcessors())); } private static ThreadPoolExecutor multiThreadedStage(Stage stage, int numThreads) diff --git a/src/java/org/apache/cassandra/db/RangeSliceCommand.java b/src/java/org/apache/cassandra/db/RangeSliceCommand.java index 5bf6f7a2eb..fd37da8b3e 100644 --- a/src/java/org/apache/cassandra/db/RangeSliceCommand.java +++ b/src/java/org/apache/cassandra/db/RangeSliceCommand.java @@ -47,7 +47,6 @@ import org.apache.cassandra.dht.AbstractBounds; import org.apache.cassandra.io.ICompactSerializer; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.net.Message; -import org.apache.cassandra.service.IReadCommand; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.thrift.ColumnParent; import org.apache.cassandra.thrift.SlicePredicate; @@ -57,7 +56,7 @@ import org.apache.thrift.TDeserializer; import org.apache.thrift.TSerializer; import org.apache.cassandra.thrift.TBinaryProtocol; -public class RangeSliceCommand implements IReadCommand +public class RangeSliceCommand { private static final RangeSliceCommandSerializer serializer = new RangeSliceCommandSerializer(); @@ -114,11 +113,6 @@ public class RangeSliceCommand implements IReadCommand ByteArrayInputStream bis = new ByteArrayInputStream(bytes); return serializer.deserialize(new DataInputStream(bis)); } - - public String getKeyspace() - { - return keyspace; - } } class RangeSliceCommandSerializer implements ICompactSerializer diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index 910b0a99db..d97489d182 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -30,12 +30,11 @@ import org.apache.cassandra.db.filter.QueryPath; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.io.ICompactSerializer; import org.apache.cassandra.net.Message; -import org.apache.cassandra.service.IReadCommand; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.FBUtilities; -public abstract class ReadCommand implements IReadCommand +public abstract class ReadCommand { public static final byte CMD_TYPE_GET_SLICE_BY_NAMES = 1; public static final byte CMD_TYPE_GET_SLICE = 2; @@ -92,11 +91,6 @@ public abstract class ReadCommand implements IReadCommand { return ColumnFamily.getComparatorFor(table, getColumnFamilyName(), queryPath.superColumnName); } - - public String getKeyspace() - { - return table; - } } class ReadCommandSerializer implements ICompactSerializer diff --git a/src/java/org/apache/cassandra/dht/BootStrapper.java b/src/java/org/apache/cassandra/dht/BootStrapper.java index cedc535698..f263dc71b0 100644 --- a/src/java/org/apache/cassandra/dht/BootStrapper.java +++ b/src/java/org/apache/cassandra/dht/BootStrapper.java @@ -282,10 +282,5 @@ public class BootStrapper token = StorageService.getPartitioner().getTokenFactory().fromString(new String(msg.getMessageBody(), Charsets.UTF_8)); condition.signalAll(); } - - public boolean isLatencyForSnitch() - { - return false; - } } } diff --git a/src/java/org/apache/cassandra/net/AsyncResult.java b/src/java/org/apache/cassandra/net/AsyncResult.java index 9f109d84e4..2819da8b1f 100644 --- a/src/java/org/apache/cassandra/net/AsyncResult.java +++ b/src/java/org/apache/cassandra/net/AsyncResult.java @@ -96,11 +96,6 @@ class AsyncResult implements IAsyncResult } } - public boolean isLatencyForSnitch() - { - return false; - } - public InetAddress getFrom() { return from; diff --git a/src/java/org/apache/cassandra/net/IMessageCallback.java b/src/java/org/apache/cassandra/net/IMessageCallback.java index d9c57eff17..96c9c566f3 100644 --- a/src/java/org/apache/cassandra/net/IMessageCallback.java +++ b/src/java/org/apache/cassandra/net/IMessageCallback.java @@ -23,9 +23,4 @@ package org.apache.cassandra.net; public interface IMessageCallback { - /** - * @return true if this callback is on the read path and its latency should be - * given as input to the dynamic snitch. - */ - public boolean isLatencyForSnitch(); } diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index fe01f327e7..2949acbc77 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -138,7 +138,7 @@ public final class MessagingService implements MessagingServiceMBean */ public void maybeAddLatency(IMessageCallback cb, InetAddress address, double latency) { - if (cb.isLatencyForSnitch()) + if (cb instanceof ReadCallback || cb instanceof AsyncResult) addLatency(address, latency); } diff --git a/src/java/org/apache/cassandra/service/AbstractRowResolver.java b/src/java/org/apache/cassandra/service/AbstractRowResolver.java deleted file mode 100644 index 7b0bdf8308..0000000000 --- a/src/java/org/apache/cassandra/service/AbstractRowResolver.java +++ /dev/null @@ -1,70 +0,0 @@ -package org.apache.cassandra.service; - -import java.io.ByteArrayInputStream; -import java.io.DataInputStream; -import java.io.IOError; -import java.io.IOException; -import java.nio.ByteBuffer; -import java.util.concurrent.ConcurrentMap; - -import org.apache.commons.lang.ArrayUtils; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import org.apache.cassandra.db.DecoratedKey; -import org.apache.cassandra.db.ReadResponse; -import org.apache.cassandra.db.Row; -import org.apache.cassandra.net.Message; -import org.apache.cassandra.utils.FBUtilities; -import org.cliffc.high_scale_lib.NonBlockingHashMap; - -public abstract class AbstractRowResolver implements IResponseResolver -{ - protected static Logger logger = LoggerFactory.getLogger(AbstractRowResolver.class); - - private static final Message FAKE_MESSAGE = new Message(FBUtilities.getLocalAddress(), StorageService.Verb.INTERNAL_RESPONSE, ArrayUtils.EMPTY_BYTE_ARRAY); - - protected final String table; - protected final ConcurrentMap replies = new NonBlockingHashMap(); - protected final DecoratedKey key; - - public AbstractRowResolver(ByteBuffer key, String table) - { - this.key = StorageService.getPartitioner().decorateKey(key); - this.table = table; - } - - public void preprocess(Message message) - { - byte[] body = message.getMessageBody(); - ByteArrayInputStream bufIn = new ByteArrayInputStream(body); - try - { - ReadResponse result = ReadResponse.serializer().deserialize(new DataInputStream(bufIn)); - if (logger.isDebugEnabled()) - logger.debug("Preprocessed {} response", result.isDigestQuery() ? "digest" : "data"); - replies.put(message, result); - } - catch (IOException e) - { - throw new IOError(e); - } - } - - /** hack so local reads don't force de/serialization of an extra real Message */ - public void injectPreProcessed(ReadResponse result) - { - assert replies.get(FAKE_MESSAGE) == null; // should only be one local reply - replies.put(FAKE_MESSAGE, result); - } - - public Iterable getMessages() - { - return replies.keySet(); - } - - public int getMessageCount() - { - return replies.size(); - } -} diff --git a/src/java/org/apache/cassandra/service/AsyncRepairCallback.java b/src/java/org/apache/cassandra/service/AsyncRepairCallback.java deleted file mode 100644 index aadbf51f29..0000000000 --- a/src/java/org/apache/cassandra/service/AsyncRepairCallback.java +++ /dev/null @@ -1,41 +0,0 @@ -package org.apache.cassandra.service; - -import java.io.IOException; - -import org.apache.cassandra.concurrent.Stage; -import org.apache.cassandra.concurrent.StageManager; -import org.apache.cassandra.net.IAsyncCallback; -import org.apache.cassandra.net.Message; -import org.apache.cassandra.utils.WrappedRunnable; - -public class AsyncRepairCallback implements IAsyncCallback -{ - private final RowRepairResolver repairResolver; - private final int count; - - public AsyncRepairCallback(RowRepairResolver repairResolver, int count) - { - this.repairResolver = repairResolver; - this.count = count; - } - - public void response(Message message) - { - repairResolver.preprocess(message); - if (repairResolver.getMessageCount() == count) - { - StageManager.getStage(Stage.READ_REPAIR).execute(new WrappedRunnable() - { - protected void runMayThrow() throws DigestMismatchException, IOException - { - repairResolver.resolve(); - } - }); - } - } - - public boolean isLatencyForSnitch() - { - return true; - } -} diff --git a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java index 9cdfa19c64..db92f6d893 100644 --- a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java +++ b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java @@ -22,12 +22,12 @@ package org.apache.cassandra.service; import java.net.InetAddress; -import java.util.List; +import java.util.Collection; import java.util.concurrent.atomic.AtomicInteger; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.db.Table; +import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.locator.IEndpointSnitch; import org.apache.cassandra.locator.NetworkTopologyStrategy; import org.apache.cassandra.net.Message; @@ -44,12 +44,12 @@ public class DatacenterReadCallback extends ReadCallback private static final String localdc = snitch.getDatacenter(FBUtilities.getLocalAddress()); private AtomicInteger localResponses; - public DatacenterReadCallback(IResponseResolver resolver, ConsistencyLevel consistencyLevel, IReadCommand command, List endpoints) + public DatacenterReadCallback(IResponseResolver resolver, ConsistencyLevel consistencyLevel, String table) { - super(resolver, consistencyLevel, command, endpoints); + super(resolver, consistencyLevel, table); localResponses = new AtomicInteger(blockfor); } - + @Override public void response(Message message) { @@ -68,15 +68,14 @@ public class DatacenterReadCallback extends ReadCallback @Override public void response(ReadResponse result) { - ((RowDigestResolver) resolver).injectPreProcessed(result); + ((ReadResponseResolver) resolver).injectPreProcessed(result); int n = localResponses.decrementAndGet(); + if (n == 0 && resolver.isDataPresent()) { condition.signal(); } - - maybeResolveForRepair(); } @Override @@ -87,7 +86,7 @@ public class DatacenterReadCallback extends ReadCallback } @Override - public void assureSufficientLiveNodes() throws UnavailableException + public void assureSufficientLiveNodes(Collection endpoints) throws UnavailableException { int localEndpoints = 0; for (InetAddress endpoint : endpoints) diff --git a/src/java/org/apache/cassandra/service/DatacenterSyncWriteResponseHandler.java b/src/java/org/apache/cassandra/service/DatacenterSyncWriteResponseHandler.java index d426a24186..1b0342b184 100644 --- a/src/java/org/apache/cassandra/service/DatacenterSyncWriteResponseHandler.java +++ b/src/java/org/apache/cassandra/service/DatacenterSyncWriteResponseHandler.java @@ -115,9 +115,4 @@ public class DatacenterSyncWriteResponseHandler extends AbstractWriteResponseHan throw new UnavailableException(); } } - - public boolean isLatencyForSnitch() - { - return false; - } } diff --git a/src/java/org/apache/cassandra/service/IReadCommand.java b/src/java/org/apache/cassandra/service/IReadCommand.java deleted file mode 100644 index 03fbc2a012..0000000000 --- a/src/java/org/apache/cassandra/service/IReadCommand.java +++ /dev/null @@ -1,6 +0,0 @@ -package org.apache.cassandra.service; - -public interface IReadCommand -{ - public String getKeyspace(); -} diff --git a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java index 4118dfe485..6660cf74e7 100644 --- a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java +++ b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java @@ -49,6 +49,7 @@ public class RangeSliceResponseResolver implements IResponseResolver> public RangeSliceResponseResolver(String table, List sources) { + assert sources.size() > 0; this.sources = sources; this.table = table; } @@ -102,8 +103,8 @@ public class RangeSliceResponseResolver implements IResponseResolver> protected Row getReduced() { - ColumnFamily resolved = RowRepairResolver.resolveSuperset(versions); - RowRepairResolver.maybeScheduleRepairs(resolved, table, key, versions, versionSources); + ColumnFamily resolved = ReadResponseResolver.resolveSuperset(versions); + ReadResponseResolver.maybeScheduleRepairs(resolved, table, key, versions, versionSources); versions.clear(); versionSources.clear(); return new Row(key, resolved); diff --git a/src/java/org/apache/cassandra/service/ReadCallback.java b/src/java/org/apache/cassandra/service/ReadCallback.java index d5587a969b..d7d8b895ac 100644 --- a/src/java/org/apache/cassandra/service/ReadCallback.java +++ b/src/java/org/apache/cassandra/service/ReadCallback.java @@ -20,20 +20,14 @@ package org.apache.cassandra.service; import java.io.IOException; import java.net.InetAddress; -import java.util.List; -import java.util.Random; +import java.util.Collection; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; -import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.concurrent.Stage; -import org.apache.cassandra.concurrent.StageManager; -import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.ReadCommand; import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.db.Table; import org.apache.cassandra.net.IAsyncCallback; @@ -42,61 +36,28 @@ import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.thrift.ConsistencyLevel; import org.apache.cassandra.thrift.UnavailableException; import org.apache.cassandra.utils.SimpleCondition; -import org.apache.cassandra.utils.WrappedRunnable; public class ReadCallback implements IAsyncCallback { protected static final Logger logger = LoggerFactory.getLogger( ReadCallback.class ); - private static final ThreadLocal random = new ThreadLocal() - { - @Override - protected Random initialValue() - { - return new Random(); - } - }; - public final IResponseResolver resolver; protected final SimpleCondition condition = new SimpleCondition(); private final long startTime; protected final int blockfor; - final List endpoints; - private final IReadCommand command; - + /** * Constructor when response count has to be calculated and blocked for. */ - public ReadCallback(IResponseResolver resolver, ConsistencyLevel consistencyLevel, IReadCommand command, List endpoints) + public ReadCallback(IResponseResolver resolver, ConsistencyLevel consistencyLevel, String table) { - this.command = command; - this.blockfor = determineBlockFor(consistencyLevel, command.getKeyspace()); + this.blockfor = determineBlockFor(consistencyLevel, table); this.resolver = resolver; this.startTime = System.currentTimeMillis(); - boolean repair = randomlyReadRepair(); - this.endpoints = repair || resolver instanceof RowRepairResolver - ? endpoints - : endpoints.subList(0, Math.min(endpoints.size(), blockfor)); // min so as to not throw exception until assureSufficient is called - if (logger.isDebugEnabled()) - logger.debug(String.format("Blockfor/repair is %s/%s; setting up requests to %s", - blockfor, repair, StringUtils.join(this.endpoints, ","))); + logger.debug("ReadCallback blocking for {} responses", blockfor); } - private boolean randomlyReadRepair() - { - if (resolver instanceof RowDigestResolver) - { - assert command instanceof ReadCommand : command; - String table = ((RowDigestResolver) resolver).table; - String columnFamily = ((ReadCommand) command).getColumnFamilyName(); - CFMetaData cfmd = DatabaseDescriptor.getTableMetaData(table).get(columnFamily); - return cfmd.getReadRepairChance() > random.get().nextDouble(); - } - // we don't read repair on range scans - return false; - } - public T get() throws TimeoutException, DigestMismatchException, IOException { long timeout = DatabaseDescriptor.getRpcTimeout() - (System.currentTimeMillis() - startTime); @@ -124,42 +85,21 @@ public class ReadCallback implements IAsyncCallback public void response(Message message) { resolver.preprocess(message); - assert resolver.getMessageCount() <= endpoints.size(); if (resolver.getMessageCount() < blockfor) return; if (resolver.isDataPresent()) - { condition.signal(); - maybeResolveForRepair(); - } } public void response(ReadResponse result) { - ((RowDigestResolver) resolver).injectPreProcessed(result); - assert resolver.getMessageCount() <= endpoints.size(); + ((ReadResponseResolver) resolver).injectPreProcessed(result); if (resolver.getMessageCount() < blockfor) return; if (resolver.isDataPresent()) - { condition.signal(); - maybeResolveForRepair(); - } } - - /** - * Check digests in the background on the Repair stage if we've received replies - * too all the requests we sent. - */ - protected void maybeResolveForRepair() - { - if (blockfor < endpoints.size() && resolver.getMessageCount() == endpoints.size()) - { - assert resolver.isDataPresent(); - StageManager.getStage(Stage.READ_REPAIR).execute(new AsyncRepairRunner()); - } - } - + public int determineBlockFor(ConsistencyLevel consistencyLevel, String table) { switch (consistencyLevel) @@ -176,38 +116,9 @@ public class ReadCallback implements IAsyncCallback } } - public void assureSufficientLiveNodes() throws UnavailableException + public void assureSufficientLiveNodes(Collection endpoints) throws UnavailableException { if (endpoints.size() < blockfor) throw new UnavailableException(); } - - public boolean isLatencyForSnitch() - { - return true; - } - - private class AsyncRepairRunner extends WrappedRunnable - { - protected void runMayThrow() throws IOException - { - try - { - resolver.resolve(); - } - catch (DigestMismatchException e) - { - if (logger.isDebugEnabled()) - logger.debug("Digest mismatch:", e); - - ReadCommand readCommand = (ReadCommand) command; - final RowRepairResolver repairResolver = new RowRepairResolver(readCommand.table, readCommand.key); - IAsyncCallback repairHandler = new AsyncRepairCallback(repairResolver, endpoints.size()); - - Message messageRepair = readCommand.makeReadMessage(); - for (InetAddress endpoint : endpoints) - MessagingService.instance().sendRR(messageRepair, endpoint, repairHandler); - } - } - } } diff --git a/src/java/org/apache/cassandra/service/ReadResponseResolver.java b/src/java/org/apache/cassandra/service/ReadResponseResolver.java index e69de29bb2..de3d1ea061 100644 --- a/src/java/org/apache/cassandra/service/ReadResponseResolver.java +++ b/src/java/org/apache/cassandra/service/ReadResponseResolver.java @@ -0,0 +1,257 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.cassandra.service; + +import java.io.ByteArrayInputStream; +import java.io.DataInputStream; +import java.io.IOError; +import java.io.IOException; +import java.net.InetAddress; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentMap; + +import org.apache.commons.lang.ArrayUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.cassandra.db.*; +import org.apache.cassandra.net.Message; +import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.utils.FBUtilities; +import org.cliffc.high_scale_lib.NonBlockingHashMap; + +/** + * Turns ReadResponse messages into Row objects, resolving to the most recent + * version and setting up read repairs as necessary. + */ +public class ReadResponseResolver implements IResponseResolver +{ + private static Logger logger_ = LoggerFactory.getLogger(ReadResponseResolver.class); + private final String table; + private final ConcurrentMap results = new NonBlockingHashMap(); + private DecoratedKey key; + private ByteBuffer digest; + private static final Message FAKE_MESSAGE = new Message(FBUtilities.getLocalAddress(), StorageService.Verb.INTERNAL_RESPONSE, ArrayUtils.EMPTY_BYTE_ARRAY);; + + public ReadResponseResolver(String table, ByteBuffer key) + { + this.table = table; + this.key = StorageService.getPartitioner().decorateKey(key); + } + + public Row getData() throws IOException + { + for (Map.Entry entry : results.entrySet()) + { + ReadResponse result = entry.getValue(); + if (!result.isDigestQuery()) + return result.row(); + } + + throw new AssertionError("getData should not be invoked when no data is present"); + } + + /* + * This method handles three different scenarios: + * + * 1a)we're handling the initial read, of data from the closest replica + digests + * from the rest. In this case we check the digests against each other, + * throw an exception if there is a mismatch, otherwise return the data row. + * + * 1b)we're checking additional digests that arrived after the minimum to handle + * the requested ConsistencyLevel, i.e. asynchronouse read repair check + * + * 2) there was a mismatch on the initial read (1a or 1b), so we redid the digest requests + * as full data reads. In this case we need to compute the most recent version + * of each column, and send diffs to out-of-date replicas. + */ + public Row resolve() throws DigestMismatchException, IOException + { + if (logger_.isDebugEnabled()) + logger_.debug("resolving " + results.size() + " responses"); + + long startTime = System.currentTimeMillis(); + List versions = new ArrayList(); + List endpoints = new ArrayList(); + + // case 1: validate digests against each other; throw immediately on mismatch. + // also, collects data results into versions/endpoints lists. + // + // results are cleared as we process them, to avoid unnecessary duplication of work + // when resolve() is called a second time for read repair on responses that were not + // necessary to satisfy ConsistencyLevel. + for (Map.Entry entry : results.entrySet()) + { + ReadResponse result = entry.getValue(); + Message message = entry.getKey(); + if (result.isDigestQuery()) + { + if (digest == null) + { + digest = result.digest(); + } + else + { + ByteBuffer digest2 = result.digest(); + if (!digest.equals(digest2)) + throw new DigestMismatchException(key, digest, digest2); + } + } + else + { + versions.add(result.row().cf); + endpoints.add(message.getFrom()); + } + + results.remove(message); + } + + // If there was a digest query compare it with all the data digests + // If there is a mismatch then throw an exception so that read repair can happen. + // + // It's important to note that we do not compare the digests of multiple data responses -- + // if we are in that situation we know there was a previous mismatch and now we're doing a repair, + // so our job is now case 2: figure out what the most recent version is and update everyone to that version. + if (digest != null) + { + for (ColumnFamily cf : versions) + { + ByteBuffer digest2 = ColumnFamily.digest(cf); + if (!digest.equals(digest2)) + throw new DigestMismatchException(key, digest, digest2); + } + if (logger_.isDebugEnabled()) + logger_.debug("digests verified"); + } + + ColumnFamily resolved; + if (versions.size() > 1) + { + resolved = resolveSuperset(versions); + if (logger_.isDebugEnabled()) + logger_.debug("versions merged"); + maybeScheduleRepairs(resolved, table, key, versions, endpoints); + } + else + { + resolved = versions.get(0); + } + + if (logger_.isDebugEnabled()) + logger_.debug("resolve: " + (System.currentTimeMillis() - startTime) + " ms."); + return new Row(key, resolved); + } + + /** + * For each row version, compare with resolved (the superset of all row versions); + * if it is missing anything, send a mutation to the endpoint it come from. + */ + public static void maybeScheduleRepairs(ColumnFamily resolved, String table, DecoratedKey key, List versions, List endpoints) + { + for (int i = 0; i < versions.size(); i++) + { + ColumnFamily diffCf = ColumnFamily.diff(versions.get(i), resolved); + if (diffCf == null) // no repair needs to happen + continue; + + // create and send the row mutation message based on the diff + RowMutation rowMutation = new RowMutation(table, key.key); + rowMutation.add(diffCf); + Message repairMessage; + try + { + repairMessage = rowMutation.makeRowMutationMessage(StorageService.Verb.READ_REPAIR); + } + catch (IOException e) + { + throw new IOError(e); + } + MessagingService.instance().sendOneWay(repairMessage, endpoints.get(i)); + } + } + + static ColumnFamily resolveSuperset(List versions) + { + assert versions.size() > 0; + + ColumnFamily resolved = null; + for (ColumnFamily cf : versions) + { + if (cf != null) + { + resolved = cf.cloneMe(); + break; + } + } + if (resolved == null) + return null; + + for (ColumnFamily cf : versions) + resolved.resolve(cf); + + return resolved; + } + + public void preprocess(Message message) + { + byte[] body = message.getMessageBody(); + ByteArrayInputStream bufIn = new ByteArrayInputStream(body); + try + { + ReadResponse result = ReadResponse.serializer().deserialize(new DataInputStream(bufIn)); + if (logger_.isDebugEnabled()) + logger_.debug("Preprocessed {} response", result.isDigestQuery() ? "digest" : "data"); + results.put(message, result); + } + catch (IOException e) + { + throw new IOError(e); + } + } + + /** hack so local reads don't force de/serialization of an extra real Message */ + public void injectPreProcessed(ReadResponse result) + { + assert results.get(FAKE_MESSAGE) == null; // should only be one local reply + results.put(FAKE_MESSAGE, result); + } + + public boolean isDataPresent() + { + for (ReadResponse result : results.values()) + { + if (!result.isDigestQuery()) + return true; + } + return false; + } + + public Iterable getMessages() + { + return results.keySet(); + } + + public int getMessageCount() + { + return results.size(); + } +} diff --git a/src/java/org/apache/cassandra/service/RepairCallback.java b/src/java/org/apache/cassandra/service/RepairCallback.java index 6242b5a8da..8ddd4849c3 100644 --- a/src/java/org/apache/cassandra/service/RepairCallback.java +++ b/src/java/org/apache/cassandra/service/RepairCallback.java @@ -39,14 +39,6 @@ public class RepairCallback implements IAsyncCallback private final SimpleCondition condition = new SimpleCondition(); private final long startTime; - /** - * The main difference between this and ReadCallback is, ReadCallback has a ConsistencyLevel - * it needs to achieve. Repair on the other hand is happy to repair whoever replies within the timeout. - * - * (The other main difference of course is, this is only created once we know we have a digest - * mismatch, and we're going to do full-data reads from everyone -- that is, this is the final - * stage in the read process.) - */ public RepairCallback(IResponseResolver resolver, List endpoints) { this.resolver = resolver; @@ -54,6 +46,10 @@ public class RepairCallback implements IAsyncCallback this.startTime = System.currentTimeMillis(); } + /** + * The main difference between this and ReadCallback is, ReadCallback has a ConsistencyLevel + * it needs to achieve. Repair on the other hand is happy to repair whoever replies within the timeout. + */ public T get() throws TimeoutException, DigestMismatchException, IOException { long timeout = DatabaseDescriptor.getRpcTimeout() - (System.currentTimeMillis() - startTime); @@ -75,9 +71,4 @@ public class RepairCallback implements IAsyncCallback if (resolver.getMessageCount() == endpoints.size()) condition.signal(); } - - public boolean isLatencyForSnitch() - { - return true; - } } diff --git a/src/java/org/apache/cassandra/service/RowDigestResolver.java b/src/java/org/apache/cassandra/service/RowDigestResolver.java deleted file mode 100644 index cd8a44a957..0000000000 --- a/src/java/org/apache/cassandra/service/RowDigestResolver.java +++ /dev/null @@ -1,125 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.cassandra.service; - -import java.io.IOException; -import java.nio.ByteBuffer; -import java.util.Map; - -import org.apache.cassandra.db.ColumnFamily; -import org.apache.cassandra.db.ReadResponse; -import org.apache.cassandra.db.Row; -import org.apache.cassandra.net.Message; - -public class RowDigestResolver extends AbstractRowResolver -{ - public RowDigestResolver(String table, ByteBuffer key) - { - super(key, table); - } - - public Row getData() throws IOException - { - for (Map.Entry entry : replies.entrySet()) - { - ReadResponse result = entry.getValue(); - if (!result.isDigestQuery()) - return result.row(); - } - - throw new AssertionError("getData should not be invoked when no data is present"); - } - - /* - * This method handles two different scenarios: - * - * 1a)we're handling the initial read, of data from the closest replica + digests - * from the rest. In this case we check the digests against each other, - * throw an exception if there is a mismatch, otherwise return the data row. - * - * 1b)we're checking additional digests that arrived after the minimum to handle - * the requested ConsistencyLevel, i.e. asynchronouse read repair check - */ - public Row resolve() throws DigestMismatchException, IOException - { - if (logger.isDebugEnabled()) - logger.debug("resolving " + replies.size() + " responses"); - - long startTime = System.currentTimeMillis(); - ColumnFamily data = null; - - // case 1: validate digests against each other; throw immediately on mismatch. - // also, collects data results into versions/endpoints lists. - // - // results are cleared as we process them, to avoid unnecessary duplication of work - // when resolve() is called a second time for read repair on responses that were not - // necessary to satisfy ConsistencyLevel. - ByteBuffer digest = null; - for (Map.Entry entry : replies.entrySet()) - { - ReadResponse response = entry.getValue(); - if (response.isDigestQuery()) - { - if (digest == null) - { - digest = response.digest(); - } - else - { - ByteBuffer digest2 = response.digest(); - if (!digest.equals(digest2)) - throw new DigestMismatchException(key, digest, digest2); - } - } - else - { - data = response.row().cf; - } - } - - // If there was a digest query compare it with all the data digests - // If there is a mismatch then throw an exception so that read repair can happen. - // - // It's important to note that we do not compare the digests of multiple data responses -- - // if we are in that situation we know there was a previous mismatch and now we're doing a repair, - // so our job is now case 2: figure out what the most recent version is and update everyone to that version. - if (digest != null) - { - ByteBuffer digest2 = ColumnFamily.digest(data); - if (!digest.equals(digest2)) - throw new DigestMismatchException(key, digest, digest2); - if (logger.isDebugEnabled()) - logger.debug("digests verified"); - } - - if (logger.isDebugEnabled()) - logger.debug("resolve: " + (System.currentTimeMillis() - startTime) + " ms."); - return new Row(key, data); - } - - public boolean isDataPresent() - { - for (ReadResponse result : replies.values()) - { - if (!result.isDigestQuery()) - return true; - } - return false; - } -} diff --git a/src/java/org/apache/cassandra/service/RowRepairResolver.java b/src/java/org/apache/cassandra/service/RowRepairResolver.java deleted file mode 100644 index 53d28807c1..0000000000 --- a/src/java/org/apache/cassandra/service/RowRepairResolver.java +++ /dev/null @@ -1,148 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.cassandra.service; - -import java.io.IOError; -import java.io.IOException; -import java.net.InetAddress; -import java.nio.ByteBuffer; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; - -import org.apache.cassandra.db.*; -import org.apache.cassandra.net.Message; -import org.apache.cassandra.net.MessagingService; - -public class RowRepairResolver extends AbstractRowResolver -{ - public RowRepairResolver(String table, ByteBuffer key) - { - super(key, table); - } - - /* - * This method handles the following scenario: - * - * there was a mismatch on the initial read (1a or 1b), so we redid the digest requests - * as full data reads. In this case we need to compute the most recent version - * of each column, and send diffs to out-of-date replicas. - */ - public Row resolve() throws DigestMismatchException, IOException - { - if (logger.isDebugEnabled()) - logger.debug("resolving " + replies.size() + " responses"); - - long startTime = System.currentTimeMillis(); - List versions = new ArrayList(); - List endpoints = new ArrayList(); - - // case 1: validate digests against each other; throw immediately on mismatch. - // also, collects data results into versions/endpoints lists. - // - // results are cleared as we process them, to avoid unnecessary duplication of work - // when resolve() is called a second time for read repair on responses that were not - // necessary to satisfy ConsistencyLevel. - for (Map.Entry entry : replies.entrySet()) - { - Message message = entry.getKey(); - ReadResponse response = entry.getValue(); - assert !response.isDigestQuery(); - versions.add(response.row().cf); - endpoints.add(message.getFrom()); - } - - ColumnFamily resolved; - if (versions.size() > 1) - { - resolved = resolveSuperset(versions); - if (logger.isDebugEnabled()) - logger.debug("versions merged"); - maybeScheduleRepairs(resolved, table, key, versions, endpoints); - } - else - { - resolved = versions.get(0); - } - - if (logger.isDebugEnabled()) - logger.debug("resolve: " + (System.currentTimeMillis() - startTime) + " ms."); - return new Row(key, resolved); - } - - /** - * For each row version, compare with resolved (the superset of all row versions); - * if it is missing anything, send a mutation to the endpoint it come from. - */ - public static void maybeScheduleRepairs(ColumnFamily resolved, String table, DecoratedKey key, List versions, List endpoints) - { - for (int i = 0; i < versions.size(); i++) - { - ColumnFamily diffCf = ColumnFamily.diff(versions.get(i), resolved); - if (diffCf == null) // no repair needs to happen - continue; - - // create and send the row mutation message based on the diff - RowMutation rowMutation = new RowMutation(table, key.key); - rowMutation.add(diffCf); - Message repairMessage; - try - { - repairMessage = rowMutation.makeRowMutationMessage(StorageService.Verb.READ_REPAIR); - } - catch (IOException e) - { - throw new IOError(e); - } - MessagingService.instance().sendOneWay(repairMessage, endpoints.get(i)); - } - } - - static ColumnFamily resolveSuperset(List versions) - { - assert versions.size() > 0; - - ColumnFamily resolved = null; - for (ColumnFamily cf : versions) - { - if (cf != null) - { - resolved = cf.cloneMe(); - break; - } - } - if (resolved == null) - return null; - - for (ColumnFamily cf : versions) - resolved.resolve(cf); - - return resolved; - } - - public Row getData() throws IOException - { - throw new UnsupportedOperationException(); - } - - public boolean isDataPresent() - { - throw new UnsupportedOperationException(); - } -} diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 0bc7cc7655..8a305aa1bc 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -36,6 +36,7 @@ import org.slf4j.LoggerFactory; import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.concurrent.StageManager; +import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.*; import org.apache.cassandra.db.filter.QueryFilter; @@ -58,6 +59,17 @@ public class StorageProxy implements StorageProxyMBean { private static final Logger logger = LoggerFactory.getLogger(StorageProxy.class); + private static ScheduledExecutorService repairExecutor = new ScheduledThreadPoolExecutor(1); // TODO JMX-enable this + + private static final ThreadLocal random = new ThreadLocal() + { + @Override + protected Random initialValue() + { + return new Random(); + } + }; + // mbean stuff private static final LatencyTracker readStats = new LatencyTracker(); private static final LatencyTracker rangeStats = new LatencyTracker(); @@ -66,8 +78,6 @@ public class StorageProxy implements StorageProxyMBean private static int maxHintWindow = DatabaseDescriptor.getMaxHintWindow(); private static final String UNREACHABLE = "UNREACHABLE"; - public static final StorageProxy instance = new StorageProxy(); - private StorageProxy() {} static { @@ -313,55 +323,66 @@ public class StorageProxy implements StorageProxyMBean private static List fetchRows(List commands, ConsistencyLevel consistency_level) throws IOException, UnavailableException, TimeoutException { List> readCallbacks = new ArrayList>(); + List> commandEndpoints = new ArrayList>(); List rows = new ArrayList(); + Set repairs = new HashSet(); // send out read requests for (ReadCommand command: commands) { assert !command.isDigestQuery(); - logger.debug("Command/ConsistencyLevel is {}/{}", command, consistency_level); List endpoints = StorageService.instance.getLiveNaturalEndpoints(command.table, command.key); DatabaseDescriptor.getEndpointSnitch().sortByProximity(FBUtilities.getLocalAddress(), endpoints); - RowDigestResolver resolver = new RowDigestResolver(command.table, command.key); - ReadCallback handler = getReadCallback(resolver, command, consistency_level, endpoints); - handler.assureSufficientLiveNodes(); - assert !handler.endpoints.isEmpty(); + ReadResponseResolver resolver = new ReadResponseResolver(command.table, command.key); + ReadCallback handler = getReadCallback(resolver, command.table, consistency_level); + handler.assureSufficientLiveNodes(endpoints); + // if we're not going to read repair, cut the endpoints list down to the ones required to satisfy ConsistencyLevel + if (randomlyReadRepair(command)) + { + if (endpoints.size() > handler.blockfor) + repairs.add(command); + } + else + { + endpoints = endpoints.subList(0, handler.blockfor); + } + // The data-request message is sent to dataPoint, the node that will actually get // the data for us. The other replicas are only sent a digest query. ReadCommand digestCommand = null; - if (handler.endpoints.size() > 1) + if (endpoints.size() > 1) { digestCommand = command.copy(); digestCommand.setDigestQuery(true); } - InetAddress dataPoint = handler.endpoints.get(0); + InetAddress dataPoint = endpoints.get(0); if (dataPoint.equals(FBUtilities.getLocalAddress())) { if (logger.isDebugEnabled()) - logger.debug("reading data locally"); + logger.debug("reading data for " + command + " locally"); StageManager.getStage(Stage.READ).execute(new LocalReadRunnable(command, handler)); } else { Message message = command.makeReadMessage(); if (logger.isDebugEnabled()) - logger.debug("reading data from " + dataPoint); + logger.debug("reading data for " + command + " from " + dataPoint); MessagingService.instance().sendRR(message, dataPoint, handler); } // We lazy-construct the digest Message object since it may not be necessary if we // are doing a local digest read, or no digest reads at all. Message digestMessage = null; - for (InetAddress digestPoint : handler.endpoints.subList(1, handler.endpoints.size())) + for (InetAddress digestPoint : endpoints.subList(1, endpoints.size())) { if (digestPoint.equals(FBUtilities.getLocalAddress())) { if (logger.isDebugEnabled()) - logger.debug("reading digest locally"); + logger.debug("reading digest for " + command + " locally"); StageManager.getStage(Stage.READ).execute(new LocalReadRunnable(digestCommand, handler)); } else @@ -369,45 +390,44 @@ public class StorageProxy implements StorageProxyMBean if (digestMessage == null) digestMessage = digestCommand.makeReadMessage(); if (logger.isDebugEnabled()) - logger.debug("reading digest for from " + digestPoint); + logger.debug("reading digest for " + command + " from " + digestPoint); MessagingService.instance().sendRR(digestMessage, digestPoint, handler); } } readCallbacks.add(handler); + commandEndpoints.add(endpoints); } // read results and make a second pass for any digest mismatches List> repairResponseHandlers = null; for (int i = 0; i < commands.size(); i++) { - ReadCallback handler = readCallbacks.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 = handler.get(); // CL.ONE is special cased here to ignore digests even if some have arrived + row = readCallback.get(); // CL.ONE is special cased here to ignore digests even if some have arrived if (row != null) rows.add(row); if (logger.isDebugEnabled()) logger.debug("Read: " + (System.currentTimeMillis() - startTime2) + " ms."); + + if (repairs.contains(command)) + repairExecutor.schedule(new RepairRunner(readCallback.resolver, command, endpoints), DatabaseDescriptor.getRpcTimeout(), TimeUnit.MILLISECONDS); } catch (DigestMismatchException ex) { if (logger.isDebugEnabled()) logger.debug("Digest mismatch:", ex); - - RowRepairResolver resolver = new RowRepairResolver(command.table, command.key); - RepairCallback repairHandler = new RepairCallback(resolver, handler.endpoints); - Message messageRepair = command.makeReadMessage(); - for (InetAddress endpoint : handler.endpoints) - MessagingService.instance().sendRR(messageRepair, endpoint, repairHandler); - + RepairCallback handler = repair(command, endpoints); if (repairResponseHandlers == null) repairResponseHandlers = new ArrayList>(); - repairResponseHandlers.add(repairHandler); + repairResponseHandlers.add(handler); } } @@ -456,13 +476,24 @@ public class StorageProxy implements StorageProxyMBean } } - static ReadCallback getReadCallback(IResponseResolver resolver, IReadCommand command, ConsistencyLevel consistencyLevel, List endpoints) + 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, command, endpoints); + return new DatacenterReadCallback(resolver, consistencyLevel, table); } - return new ReadCallback(resolver, consistencyLevel, command, endpoints); + return new ReadCallback(resolver, consistencyLevel, table); + } + + private static RepairCallback repair(ReadCommand command, List endpoints) + throws IOException + { + ReadResponseResolver resolver = new ReadResponseResolver(command.table, command.key); + RepairCallback handler = new RepairCallback(resolver, endpoints); + Message messageRepair = command.makeReadMessage(); + for (InetAddress endpoint : endpoints) + MessagingService.instance().sendRR(messageRepair, endpoint, handler); + return handler; } /* @@ -514,14 +545,16 @@ public class StorageProxy implements StorageProxyMBean // collect replies and resolve according to consistency level RangeSliceResponseResolver resolver = new RangeSliceResponseResolver(command.keyspace, liveEndpoints); - ReadCallback> handler = getReadCallback(resolver, command, consistency_level, liveEndpoints); - handler.assureSufficientLiveNodes(); + AbstractReplicationStrategy rs = Table.open(command.keyspace).getReplicationStrategy(); + ReadCallback> handler = getReadCallback(resolver, command.keyspace, consistency_level); + // TODO bail early if live endpoints can't satisfy requested consistency level for (InetAddress endpoint : liveEndpoints) { MessagingService.instance().sendRR(message, endpoint, handler); if (logger.isDebugEnabled()) logger.debug("reading " + c2 + " from " + endpoint); } + // TODO read repair on remaining replicas? // if we're done, great, otherwise, move to the next range try @@ -574,11 +607,6 @@ public class StorageProxy implements StorageProxyMBean versions.put(message.getFrom(), theirVersion); latch.countDown(); } - - public boolean isLatencyForSnitch() - { - return false; - } }; // an empty message acts as a request to the SchemaCheckVerbHandler. for (InetAddress endpoint : liveHosts) @@ -671,6 +699,12 @@ public class StorageProxy implements StorageProxyMBean return ranges; } + + private static boolean randomlyReadRepair(ReadCommand command) + { + CFMetaData cfmd = DatabaseDescriptor.getTableMetaData(command.table).get(command.getColumnFamilyName()); + return cfmd.getReadRepairChance() > random.get().nextDouble(); + } public long getReadOperations() { @@ -747,7 +781,7 @@ public class StorageProxy implements StorageProxyMBean return writeStats.getRecentLatencyHistogramMicros(); } - public static List scan(final String keyspace, String column_family, IndexClause index_clause, SlicePredicate column_predicate, ConsistencyLevel consistency_level) + public static List scan(String keyspace, String column_family, IndexClause index_clause, SlicePredicate column_predicate, ConsistencyLevel consistency_level) throws IOException, TimeoutException, UnavailableException { IPartitioner p = StorageService.getPartitioner(); @@ -765,16 +799,12 @@ public class StorageProxy implements StorageProxyMBean // collect replies and resolve according to consistency level RangeSliceResponseResolver resolver = new RangeSliceResponseResolver(keyspace, liveEndpoints); - IReadCommand iCommand = new IReadCommand() - { - public String getKeyspace() - { - return keyspace; - } - }; - ReadCallback> handler = getReadCallback(resolver, iCommand, consistency_level, liveEndpoints); - handler.assureSufficientLiveNodes(); - + ReadCallback> handler = getReadCallback(resolver, keyspace, consistency_level); + + // bail early if live endpoints can't satisfy requested consistency level + if(handler.blockfor > liveEndpoints.size()) + throw new UnavailableException(); + IndexScanCommand command = new IndexScanCommand(keyspace, column_family, index_clause, column_predicate, range); Message message = command.getMessage(); for (InetAddress endpoint : liveEndpoints) @@ -882,4 +912,40 @@ public class StorageProxy implements StorageProxyMBean { return !Gossiper.instance.getUnreachableMembers().isEmpty(); } + + private static class RepairRunner extends WrappedRunnable + { + private final IResponseResolver resolver; + private final ReadCommand command; + private final List endpoints; + + public RepairRunner(IResponseResolver resolver, ReadCommand command, List endpoints) + { + this.resolver = resolver; + this.command = command; + this.endpoints = endpoints; + } + + protected void runMayThrow() throws IOException + { + try + { + resolver.resolve(); + } + catch (DigestMismatchException e) + { + if (logger.isDebugEnabled()) + logger.debug("Digest mismatch:", e); + final RepairCallback callback = repair(command, endpoints); + Runnable runnable = new WrappedRunnable() + { + public void runMayThrow() throws DigestMismatchException, IOException, TimeoutException + { + callback.get(); + } + }; + repairExecutor.schedule(runnable, DatabaseDescriptor.getRpcTimeout(), TimeUnit.MILLISECONDS); + } + } + } } diff --git a/src/java/org/apache/cassandra/service/TruncateResponseHandler.java b/src/java/org/apache/cassandra/service/TruncateResponseHandler.java index 7aea6f4c74..20764bd6b5 100644 --- a/src/java/org/apache/cassandra/service/TruncateResponseHandler.java +++ b/src/java/org/apache/cassandra/service/TruncateResponseHandler.java @@ -73,9 +73,4 @@ public class TruncateResponseHandler implements IAsyncCallback if (responses.get() >= responseCount) condition.signal(); } - - public boolean isLatencyForSnitch() - { - return false; - } } diff --git a/src/java/org/apache/cassandra/service/WriteResponseHandler.java b/src/java/org/apache/cassandra/service/WriteResponseHandler.java index d5a903de40..94578686e1 100644 --- a/src/java/org/apache/cassandra/service/WriteResponseHandler.java +++ b/src/java/org/apache/cassandra/service/WriteResponseHandler.java @@ -121,9 +121,4 @@ public class WriteResponseHandler extends AbstractWriteResponseHandler throw new UnavailableException(); } } - - public boolean isLatencyForSnitch() - { - return false; - } } diff --git a/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java b/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java index 8d405b80b1..f142231664 100644 --- a/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java +++ b/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java @@ -71,7 +71,7 @@ public class ConsistencyLevelTest extends CleanupHelper AbstractReplicationStrategy strategy; - for (final String table : DatabaseDescriptor.getNonSystemTables()) + for (String table : DatabaseDescriptor.getNonSystemTables()) { strategy = getStrategy(table, tmd); StorageService.calculatePendingRanges(strategy, table); @@ -96,15 +96,7 @@ public class ConsistencyLevelTest extends CleanupHelper IWriteResponseHandler writeHandler = strategy.getWriteResponseHandler(hosts, hintedNodes, c); - IReadCommand command = new IReadCommand() - { - public String getKeyspace() - { - return table; - } - }; - RowRepairResolver resolver = new RowRepairResolver(table, ByteBufferUtil.bytes("foo")); - ReadCallback readHandler = StorageProxy.getReadCallback(resolver, command, c, new ArrayList(hintedNodes.keySet())); + ReadCallback readHandler = StorageProxy.getReadCallback(new ReadResponseResolver(table, ByteBufferUtil.bytes("foo")), table, c); boolean isWriteUnavailable = false; boolean isReadUnavailable = false; @@ -119,7 +111,7 @@ public class ConsistencyLevelTest extends CleanupHelper try { - readHandler.assureSufficientLiveNodes(); + readHandler.assureSufficientLiveNodes(hintedNodes.asMap().keySet()); } catch (UnavailableException e) { diff --git a/test/unit/org/apache/cassandra/service/RowResolverTest.java b/test/unit/org/apache/cassandra/service/ReadResponseResolverTest.java similarity index 84% rename from test/unit/org/apache/cassandra/service/RowResolverTest.java rename to test/unit/org/apache/cassandra/service/ReadResponseResolverTest.java index 4c847cc4de..8c988a2e9f 100644 --- a/test/unit/org/apache/cassandra/service/RowResolverTest.java +++ b/test/unit/org/apache/cassandra/service/ReadResponseResolverTest.java @@ -32,7 +32,7 @@ import static org.apache.cassandra.db.TableTest.assertColumns; import static org.apache.cassandra.Util.column; import static junit.framework.Assert.assertNull; -public class RowResolverTest extends SchemaLoader +public class ReadResponseResolverTest extends SchemaLoader { @Test public void testResolveSupersetNewer() @@ -43,7 +43,7 @@ public class RowResolverTest extends SchemaLoader ColumnFamily cf2 = ColumnFamily.create("Keyspace1", "Standard1"); cf2.addColumn(column("c1", "v2", 1)); - ColumnFamily resolved = RowRepairResolver.resolveSuperset(Arrays.asList(cf1, cf2)); + ColumnFamily resolved = ReadResponseResolver.resolveSuperset(Arrays.asList(cf1, cf2)); assertColumns(resolved, "c1"); assertColumns(ColumnFamily.diff(cf1, resolved), "c1"); assertNull(ColumnFamily.diff(cf2, resolved)); @@ -58,7 +58,7 @@ public class RowResolverTest extends SchemaLoader ColumnFamily cf2 = ColumnFamily.create("Keyspace1", "Standard1"); cf2.addColumn(column("c2", "v2", 1)); - ColumnFamily resolved = RowRepairResolver.resolveSuperset(Arrays.asList(cf1, cf2)); + ColumnFamily resolved = ReadResponseResolver.resolveSuperset(Arrays.asList(cf1, cf2)); assertColumns(resolved, "c1", "c2"); assertColumns(ColumnFamily.diff(cf1, resolved), "c2"); assertColumns(ColumnFamily.diff(cf2, resolved), "c1"); @@ -70,7 +70,7 @@ public class RowResolverTest extends SchemaLoader ColumnFamily cf2 = ColumnFamily.create("Keyspace1", "Standard1"); cf2.addColumn(column("c2", "v2", 1)); - ColumnFamily resolved = RowRepairResolver.resolveSuperset(Arrays.asList(null, cf2)); + ColumnFamily resolved = ReadResponseResolver.resolveSuperset(Arrays.asList(null, cf2)); assertColumns(resolved, "c2"); assertColumns(ColumnFamily.diff(null, resolved), "c2"); assertNull(ColumnFamily.diff(cf2, resolved)); @@ -82,7 +82,7 @@ public class RowResolverTest extends SchemaLoader ColumnFamily cf1 = ColumnFamily.create("Keyspace1", "Standard1"); cf1.addColumn(column("c1", "v1", 0)); - ColumnFamily resolved = RowRepairResolver.resolveSuperset(Arrays.asList(cf1, null)); + ColumnFamily resolved = ReadResponseResolver.resolveSuperset(Arrays.asList(cf1, null)); assertColumns(resolved, "c1"); assertNull(ColumnFamily.diff(cf1, resolved)); assertColumns(ColumnFamily.diff(null, resolved), "c1"); @@ -91,6 +91,6 @@ public class RowResolverTest extends SchemaLoader @Test public void testResolveSupersetNullBoth() { - assertNull(RowRepairResolver.resolveSuperset(Arrays.asList(null, null))); + assertNull(ReadResponseResolver.resolveSuperset(Arrays.asList(null, null))); } }