From f5e1cbca871e6e4cc7007177b2b0e9f367ae60ba Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Wed, 19 Mar 2014 00:15:19 -0500 Subject: [PATCH] Include correct consistencyLevel in LWT timeout patch by Sankalp Kohli; reviewed by jbellis for CASSANDRA-6884 --- CHANGES.txt | 1 + .../apache/cassandra/service/StorageProxy.java | 16 ++++++++-------- .../service/paxos/AbstractPaxosCallback.java | 6 ++++-- .../cassandra/service/paxos/PrepareCallback.java | 5 +++-- .../cassandra/service/paxos/ProposeCallback.java | 5 +++-- 5 files changed, 19 insertions(+), 14 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 4741475533..7eebd5b130 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.0.7 + * Include correct consistencyLevel in LWT timeout (CASSANDRA-6884) * Lower chances for losing new SSTables during nodetool refresh and ColumnFamilyStore.loadNewSSTables (CASSANDRA-6514) * Add support for DELETE ... IF EXISTS to CQL3 (CASSANDRA-5708) diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index fda9819a7e..a6912c2341 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -246,7 +246,7 @@ public class StorageProxy implements StorageProxyMBean Commit proposal = Commit.newProposal(key, ballot, updates); Tracing.trace("CAS precondition is met; proposing client-requested updates for {}", ballot); - if (proposePaxos(proposal, liveEndpoints, requiredParticipants, true)) + if (proposePaxos(proposal, liveEndpoints, requiredParticipants, true, consistencyForPaxos)) { if (consistencyForCommit == ConsistencyLevel.ANY) sendCommit(proposal, liveEndpoints); @@ -318,7 +318,7 @@ public class StorageProxy implements StorageProxyMBean // prepare Tracing.trace("Preparing {}", ballot); Commit toPrepare = Commit.newPrepare(key, metadata, ballot); - summary = preparePaxos(toPrepare, liveEndpoints, requiredParticipants); + summary = preparePaxos(toPrepare, liveEndpoints, requiredParticipants, consistencyForPaxos); if (!summary.promised) { Tracing.trace("Some replicas have already promised a higher ballot than ours; aborting"); @@ -336,7 +336,7 @@ public class StorageProxy implements StorageProxyMBean { Tracing.trace("Finishing incomplete paxos round {}", inProgress); Commit refreshedInProgress = Commit.newProposal(inProgress.key, ballot, inProgress.update); - if (proposePaxos(refreshedInProgress, liveEndpoints, requiredParticipants, false)) + if (proposePaxos(refreshedInProgress, liveEndpoints, requiredParticipants, false, consistencyForPaxos)) { commitPaxos(refreshedInProgress, ConsistencyLevel.QUORUM); } @@ -381,10 +381,10 @@ public class StorageProxy implements StorageProxyMBean MessagingService.instance().sendOneWay(message, target); } - private static PrepareCallback preparePaxos(Commit toPrepare, List endpoints, int requiredParticipants) + private static PrepareCallback preparePaxos(Commit toPrepare, List endpoints, int requiredParticipants, ConsistencyLevel consistencyForPaxos) throws WriteTimeoutException { - PrepareCallback callback = new PrepareCallback(toPrepare.key, toPrepare.update.metadata(), requiredParticipants); + PrepareCallback callback = new PrepareCallback(toPrepare.key, toPrepare.update.metadata(), requiredParticipants, consistencyForPaxos); MessageOut message = new MessageOut(MessagingService.Verb.PAXOS_PREPARE, toPrepare, Commit.serializer); for (InetAddress target : endpoints) MessagingService.instance().sendRR(message, target, callback); @@ -392,10 +392,10 @@ public class StorageProxy implements StorageProxyMBean return callback; } - private static boolean proposePaxos(Commit proposal, List endpoints, int requiredParticipants, boolean timeoutIfPartial) + private static boolean proposePaxos(Commit proposal, List endpoints, int requiredParticipants, boolean timeoutIfPartial, ConsistencyLevel consistencyLevel) throws WriteTimeoutException { - ProposeCallback callback = new ProposeCallback(endpoints.size(), requiredParticipants, !timeoutIfPartial); + ProposeCallback callback = new ProposeCallback(endpoints.size(), requiredParticipants, !timeoutIfPartial, consistencyLevel); MessageOut message = new MessageOut(MessagingService.Verb.PAXOS_PROPOSE, proposal, Commit.serializer); for (InetAddress target : endpoints) MessagingService.instance().sendRR(message, target, callback); @@ -406,7 +406,7 @@ public class StorageProxy implements StorageProxyMBean return true; if (timeoutIfPartial && !callback.isFullyRefused()) - throw new WriteTimeoutException(WriteType.CAS, ConsistencyLevel.SERIAL, callback.getAcceptCount(), requiredParticipants); + throw new WriteTimeoutException(WriteType.CAS, consistencyLevel, callback.getAcceptCount(), requiredParticipants); return false; } diff --git a/src/java/org/apache/cassandra/service/paxos/AbstractPaxosCallback.java b/src/java/org/apache/cassandra/service/paxos/AbstractPaxosCallback.java index 8197cfd61b..37defde2ad 100644 --- a/src/java/org/apache/cassandra/service/paxos/AbstractPaxosCallback.java +++ b/src/java/org/apache/cassandra/service/paxos/AbstractPaxosCallback.java @@ -34,10 +34,12 @@ public abstract class AbstractPaxosCallback implements IAsyncCallback { protected final CountDownLatch latch; protected final int targets; + private final ConsistencyLevel consistency; - public AbstractPaxosCallback(int targets) + public AbstractPaxosCallback(int targets, ConsistencyLevel consistency) { this.targets = targets; + this.consistency = consistency; latch = new CountDownLatch(targets); } @@ -56,7 +58,7 @@ public abstract class AbstractPaxosCallback implements IAsyncCallback try { if (!latch.await(DatabaseDescriptor.getWriteRpcTimeout(), TimeUnit.MILLISECONDS)) - throw new WriteTimeoutException(WriteType.CAS, ConsistencyLevel.SERIAL, getResponseCount(), targets); + throw new WriteTimeoutException(WriteType.CAS, consistency, getResponseCount(), targets); } catch (InterruptedException ex) { diff --git a/src/java/org/apache/cassandra/service/paxos/PrepareCallback.java b/src/java/org/apache/cassandra/service/paxos/PrepareCallback.java index 04a18b93d3..a446b0b6c6 100644 --- a/src/java/org/apache/cassandra/service/paxos/PrepareCallback.java +++ b/src/java/org/apache/cassandra/service/paxos/PrepareCallback.java @@ -28,6 +28,7 @@ import java.util.concurrent.ConcurrentHashMap; import com.google.common.base.Predicate; import com.google.common.collect.Iterables; +import org.apache.cassandra.db.ConsistencyLevel; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -45,9 +46,9 @@ public class PrepareCallback extends AbstractPaxosCallback private final Map commitsByReplica = new ConcurrentHashMap(); - public PrepareCallback(ByteBuffer key, CFMetaData metadata, int targets) + public PrepareCallback(ByteBuffer key, CFMetaData metadata, int targets, ConsistencyLevel consistency) { - super(targets); + super(targets, consistency); // need to inject the right key in the empty commit so comparing with empty commits in the reply works as expected mostRecentCommit = Commit.emptyCommit(key, metadata); mostRecentInProgressCommit = Commit.emptyCommit(key, metadata); diff --git a/src/java/org/apache/cassandra/service/paxos/ProposeCallback.java b/src/java/org/apache/cassandra/service/paxos/ProposeCallback.java index 0075840cb0..018dab91ac 100644 --- a/src/java/org/apache/cassandra/service/paxos/ProposeCallback.java +++ b/src/java/org/apache/cassandra/service/paxos/ProposeCallback.java @@ -23,6 +23,7 @@ package org.apache.cassandra.service.paxos; import java.util.concurrent.atomic.AtomicInteger; +import org.apache.cassandra.db.ConsistencyLevel; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -49,9 +50,9 @@ public class ProposeCallback extends AbstractPaxosCallback private final int requiredAccepts; private final boolean failFast; - public ProposeCallback(int totalTargets, int requiredTargets, boolean failFast) + public ProposeCallback(int totalTargets, int requiredTargets, boolean failFast, ConsistencyLevel consistency) { - super(totalTargets); + super(totalTargets, consistency); this.requiredAccepts = requiredTargets; this.failFast = failFast; }