jobs = this.jobs;
if (jobs != null)
{
- for (RepairJob job : jobs)
+ for (AbstractRepairJob job : jobs)
job.abort(reason);
}
this.jobs = null;
diff --git a/src/java/org/apache/cassandra/repair/messages/RepairMessage.java b/src/java/org/apache/cassandra/repair/messages/RepairMessage.java
index e38a930bcc..2eab88ef91 100644
--- a/src/java/org/apache/cassandra/repair/messages/RepairMessage.java
+++ b/src/java/org/apache/cassandra/repair/messages/RepairMessage.java
@@ -25,7 +25,6 @@ import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import java.util.function.Supplier;
-
import javax.annotation.Nullable;
import com.google.common.annotations.VisibleForTesting;
@@ -74,7 +73,7 @@ public abstract class RepairMessage
}
@Override
- public void onFailure(InetAddressAndPort from, RequestFailure failureReason)
+ public void onFailure(InetAddressAndPort from, RequestFailure failure)
{
}
};
diff --git a/src/java/org/apache/cassandra/repair/messages/RepairOption.java b/src/java/org/apache/cassandra/repair/messages/RepairOption.java
index bc9231dcc1..bef7acfe16 100644
--- a/src/java/org/apache/cassandra/repair/messages/RepairOption.java
+++ b/src/java/org/apache/cassandra/repair/messages/RepairOption.java
@@ -17,7 +17,13 @@
*/
package org.apache.cassandra.repair.messages;
-import java.util.*;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Set;
+import java.util.StringTokenizer;
import com.google.common.base.Joiner;
import com.google.common.base.Preconditions;
@@ -57,6 +63,8 @@ public class RepairOption
public static final String NO_TOMBSTONE_PURGING = "nopurge";
+ public static final String ACCORD_REPAIR_KEY = "accordRepair";
+
// we don't want to push nodes too much for repair
public static final int MAX_JOB_THREADS = 4;
@@ -86,6 +94,7 @@ public class RepairOption
}
return ranges;
}
+
/**
* Construct RepairOptions object from given map of Strings.
*
@@ -167,6 +176,12 @@ public class RepairOption
* ranges to the same host multiple times
*
false |
*
+ *
+ * | accordRepair |
+ * "true" if the repair should be of Accord in flight transactions. Will ensure
+ * that once repair completes all Accord transactions are replicated at quorum |
+ * false |
+ *
*
*
*
@@ -188,11 +203,21 @@ public class RepairOption
boolean repairPaxos = Boolean.parseBoolean(options.get(REPAIR_PAXOS_KEY));
boolean paxosOnly = Boolean.parseBoolean(options.get(PAXOS_ONLY_KEY));
boolean dontPurgeTombstones = Boolean.parseBoolean(options.get(NO_TOMBSTONE_PURGING));
+ boolean accordRepair = Boolean.parseBoolean(options.get(ACCORD_REPAIR_KEY));
if (previewKind != PreviewKind.NONE)
{
Preconditions.checkArgument(!repairPaxos, "repairPaxos must be set to false for preview repairs");
Preconditions.checkArgument(!paxosOnly, "paxosOnly must be set to false for preview repairs");
+ Preconditions.checkArgument(!accordRepair, "accordRepair must be set to false for preview repairs");
+ }
+
+ if (accordRepair)
+ {
+ Preconditions.checkArgument(!paxosOnly, "paxosOnly must be set to false for Accord repairs");
+ Preconditions.checkArgument(previewKind == PreviewKind.NONE, "Can't perform preview repair with an Accord repair");
+ Preconditions.checkArgument(!force, "Accord repair only requires a quorum to work so force is not supported");
+ incremental = false;
}
int jobThreads = 1;
@@ -212,7 +237,7 @@ public class RepairOption
boolean asymmetricSyncing = Boolean.parseBoolean(options.get(OPTIMISE_STREAMS_KEY));
- RepairOption option = new RepairOption(parallelism, primaryRange, incremental, trace, jobThreads, ranges, !ranges.isEmpty(), pullRepair, force, previewKind, asymmetricSyncing, ignoreUnreplicatedKeyspaces, repairPaxos, paxosOnly, dontPurgeTombstones);
+ RepairOption option = new RepairOption(parallelism, primaryRange, incremental, trace, jobThreads, ranges, pullRepair, force, previewKind, asymmetricSyncing, ignoreUnreplicatedKeyspaces, repairPaxos, paxosOnly, dontPurgeTombstones, accordRepair);
// data centers
String dataCentersStr = options.get(DATACENTERS_KEY);
@@ -286,7 +311,6 @@ public class RepairOption
private final boolean incremental;
private final boolean trace;
private final int jobThreads;
- private final boolean isSubrangeRepair;
private final boolean pullRepair;
private final boolean forceRepair;
private final PreviewKind previewKind;
@@ -296,12 +320,17 @@ public class RepairOption
private final boolean paxosOnly;
private final boolean dontPurgeTombstones;
+ private final boolean accordRepair;
+
private final Collection columnFamilies = new HashSet<>();
private final Collection dataCenters = new HashSet<>();
private final Collection hosts = new HashSet<>();
private final Collection> ranges = new HashSet<>();
- public RepairOption(RepairParallelism parallelism, boolean primaryRange, boolean incremental, boolean trace, int jobThreads, Collection> ranges, boolean isSubrangeRepair, boolean pullRepair, boolean forceRepair, PreviewKind previewKind, boolean optimiseStreams, boolean ignoreUnreplicatedKeyspaces, boolean repairPaxos, boolean paxosOnly, boolean dontPurgeTombstones)
+ public RepairOption(RepairParallelism parallelism, boolean primaryRange, boolean incremental, boolean trace, int jobThreads,
+ Collection> ranges, boolean pullRepair, boolean forceRepair,
+ PreviewKind previewKind, boolean optimiseStreams, boolean ignoreUnreplicatedKeyspaces, boolean repairPaxos,
+ boolean paxosOnly, boolean dontPurgeTombstones, boolean accordRepair)
{
this.parallelism = parallelism;
@@ -310,7 +339,6 @@ public class RepairOption
this.trace = trace;
this.jobThreads = jobThreads;
this.ranges.addAll(ranges);
- this.isSubrangeRepair = isSubrangeRepair;
this.pullRepair = pullRepair;
this.forceRepair = forceRepair;
this.previewKind = previewKind;
@@ -319,6 +347,7 @@ public class RepairOption
this.repairPaxos = repairPaxos;
this.paxosOnly = paxosOnly;
this.dontPurgeTombstones = dontPurgeTombstones;
+ this.accordRepair = accordRepair;
}
public RepairParallelism getParallelism()
@@ -381,11 +410,6 @@ public class RepairOption
return dataCenters.isEmpty() && hosts.isEmpty();
}
- public boolean isSubrangeRepair()
- {
- return isSubrangeRepair;
- }
-
public PreviewKind getPreviewKind()
{
return previewKind;
@@ -439,6 +463,11 @@ public class RepairOption
return dontPurgeTombstones;
}
+ public boolean accordRepair()
+ {
+ return accordRepair;
+ }
+
@Override
public String toString()
{
@@ -459,6 +488,7 @@ public class RepairOption
", repairPaxos: " + repairPaxos +
", paxosOnly: " + paxosOnly +
", dontPurgeTombstones: " + dontPurgeTombstones +
+ ", accordRepair: " + accordRepair +
')';
}
@@ -472,7 +502,6 @@ public class RepairOption
options.put(COLUMNFAMILIES_KEY, Joiner.on(",").join(columnFamilies));
options.put(DATACENTERS_KEY, Joiner.on(",").join(dataCenters));
options.put(HOSTS_KEY, Joiner.on(",").join(hosts));
- options.put(SUB_RANGE_REPAIR_KEY, Boolean.toString(isSubrangeRepair));
options.put(TRACE_KEY, Boolean.toString(trace));
options.put(RANGES_KEY, Joiner.on(",").join(ranges));
options.put(PULL_REPAIR_KEY, Boolean.toString(pullRepair));
@@ -482,6 +511,7 @@ public class RepairOption
options.put(REPAIR_PAXOS_KEY, Boolean.toString(repairPaxos));
options.put(PAXOS_ONLY_KEY, Boolean.toString(paxosOnly));
options.put(NO_TOMBSTONE_PURGING, Boolean.toString(dontPurgeTombstones));
+ options.put(ACCORD_REPAIR_KEY, Boolean.toString(accordRepair));
return options;
}
}
diff --git a/src/java/org/apache/cassandra/repair/messages/SyncResponse.java b/src/java/org/apache/cassandra/repair/messages/SyncResponse.java
index e7e7985fff..0c528a3796 100644
--- a/src/java/org/apache/cassandra/repair/messages/SyncResponse.java
+++ b/src/java/org/apache/cassandra/repair/messages/SyncResponse.java
@@ -23,12 +23,13 @@ import java.util.List;
import java.util.Objects;
import org.apache.cassandra.db.TypeSizes;
+import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.locator.InetAddressAndPort;
-import org.apache.cassandra.repair.SyncNodePair;
import org.apache.cassandra.repair.RepairJobDesc;
+import org.apache.cassandra.repair.SyncNodePair;
import org.apache.cassandra.streaming.SessionSummary;
/**
@@ -103,7 +104,7 @@ public class SyncResponse extends RepairMessage
List summaries = new ArrayList<>(numSummaries);
for (int i=0; i
{
return id.compareTo(o.id);
}
+
+ public static final IVersionedSerializer serializer = new IVersionedSerializer()
+ {
+ @Override
+ public void serialize(TableId t, DataOutputPlus out, int version) throws IOException
+ {
+ t.serialize(out);
+ }
+
+ @Override
+ public TableId deserialize(DataInputPlus in, int version) throws IOException
+ {
+ return TableId.deserialize(in);
+ }
+
+ @Override
+ public long serializedSize(TableId t, int version)
+ {
+ return t.serializedSize();
+ }
+ };
+
+ public static final MetadataSerializer metadataSerializer = new MetadataSerializer()
+ {
+ @Override
+ public void serialize(TableId t, DataOutputPlus out, Version version) throws IOException
+ {
+ t.serialize(out);
+ }
+
+ @Override
+ public TableId deserialize(DataInputPlus in, Version version) throws IOException
+ {
+ return TableId.deserialize(in);
+ }
+
+ @Override
+ public long serializedSize(TableId t, Version version)
+ {
+ return t.serializedSize();
+ }
+ };
}
diff --git a/src/java/org/apache/cassandra/schema/TableMetadata.java b/src/java/org/apache/cassandra/schema/TableMetadata.java
index 1c22d955b7..90c56e47bb 100644
--- a/src/java/org/apache/cassandra/schema/TableMetadata.java
+++ b/src/java/org/apache/cassandra/schema/TableMetadata.java
@@ -70,10 +70,10 @@ import org.apache.cassandra.exceptions.ConfigurationException;
import org.apache.cassandra.exceptions.InvalidRequestException;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
+import org.apache.cassandra.service.reads.SpeculativeRetryPolicy;
import org.apache.cassandra.tcm.Epoch;
import org.apache.cassandra.tcm.serialization.UDTAndFunctionsAwareMetadataSerializer;
import org.apache.cassandra.tcm.serialization.Version;
-import org.apache.cassandra.service.reads.SpeculativeRetryPolicy;
import org.apache.cassandra.utils.AbstractIterator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.FBUtilities;
@@ -320,6 +320,11 @@ public class TableMetadata implements SchemaElement
return unbuild().indexes(indexes).build();
}
+ public TableId id()
+ {
+ return id;
+ }
+
public boolean isView()
{
return kind == Kind.VIEW;
@@ -344,7 +349,7 @@ public class TableMetadata implements SchemaElement
{
return false;
}
-
+
public boolean isIncrementalBackupsEnabled()
{
return params.incrementalBackups;
diff --git a/src/java/org/apache/cassandra/service/ActiveRepairService.java b/src/java/org/apache/cassandra/service/ActiveRepairService.java
index 1592718bf1..24b2966bdf 100644
--- a/src/java/org/apache/cassandra/service/ActiveRepairService.java
+++ b/src/java/org/apache/cassandra/service/ActiveRepairService.java
@@ -448,6 +448,7 @@ public class ActiveRepairService implements IEndpointStateChangeSubscriber, IFai
*/
public RepairSession submitRepairSession(TimeUUID parentRepairSession,
CommonRange range,
+ boolean excludedDeadNodes,
String keyspace,
RepairParallelism parallelismDegree,
boolean isIncremental,
@@ -457,6 +458,7 @@ public class ActiveRepairService implements IEndpointStateChangeSubscriber, IFai
boolean repairPaxos,
boolean paxosOnly,
boolean dontPurgeTombstones,
+ boolean accordRepair,
ExecutorPlus executor,
Scheduler validationScheduler,
String... cfnames)
@@ -470,9 +472,11 @@ public class ActiveRepairService implements IEndpointStateChangeSubscriber, IFai
if (cfnames.length == 0)
return null;
- final RepairSession session = new RepairSession(ctx, validationScheduler, parentRepairSession, range, keyspace,
+ final RepairSession session = new RepairSession(ctx, validationScheduler, parentRepairSession,
+ range, excludedDeadNodes, keyspace,
parallelismDegree, isIncremental, pullRepair,
- previewKind, optimiseStreams, repairPaxos, paxosOnly, dontPurgeTombstones, cfnames);
+ previewKind, optimiseStreams, repairPaxos, paxosOnly,
+ dontPurgeTombstones, accordRepair, cfnames);
repairs.getIfPresent(parentRepairSession).register(session.state);
sessions.put(session.getId(), session);
diff --git a/src/java/org/apache/cassandra/service/CASRequest.java b/src/java/org/apache/cassandra/service/CASRequest.java
index f118dcf847..fb78daa2a5 100644
--- a/src/java/org/apache/cassandra/service/CASRequest.java
+++ b/src/java/org/apache/cassandra/service/CASRequest.java
@@ -18,15 +18,17 @@
package org.apache.cassandra.service;
import accord.primitives.Txn;
+import org.apache.cassandra.db.ConsistencyLevel;
import org.apache.cassandra.db.SinglePartitionReadCommand;
import org.apache.cassandra.db.partitions.FilteredPartition;
import org.apache.cassandra.db.partitions.PartitionUpdate;
-import org.apache.cassandra.db.rows.RowIterator;
import org.apache.cassandra.exceptions.InvalidRequestException;
-import org.apache.cassandra.service.accord.txn.TxnData;
+import org.apache.cassandra.service.accord.txn.TxnResult;
import org.apache.cassandra.service.paxos.Ballot;
import org.apache.cassandra.transport.Dispatcher;
+import static org.apache.cassandra.service.StorageProxy.ConsensusAttemptResult;
+
/**
* Abstract the conditions and updates for a CAS operation.
*/
@@ -51,7 +53,7 @@ public interface CASRequest
*/
PartitionUpdate makeUpdates(FilteredPartition current, ClientState clientState, Ballot ballot) throws InvalidRequestException;
- Txn toAccordTxn(ClientState clientState, long nowInSecs);
+ Txn toAccordTxn(ConsistencyLevel consistencyLevel, ConsistencyLevel commitConsistencyLevel, ClientState clientState, long nowInSecs);
- RowIterator toCasResult(TxnData data);
+ ConsensusAttemptResult toCasResult(TxnResult txnResult);
}
diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java
index e61c9187a3..a40426c001 100644
--- a/src/java/org/apache/cassandra/service/StorageProxy.java
+++ b/src/java/org/apache/cassandra/service/StorageProxy.java
@@ -39,14 +39,16 @@ import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Function;
import java.util.stream.Collectors;
+import javax.annotation.Nonnull;
+import javax.annotation.Nullable;
-import com.google.common.base.Preconditions;
import com.google.common.cache.CacheLoader;
import com.google.common.collect.Iterables;
import com.google.common.util.concurrent.Uninterruptibles;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import accord.primitives.Keys;
import accord.primitives.Txn;
import org.apache.cassandra.batchlog.Batch;
import org.apache.cassandra.batchlog.BatchlogManager;
@@ -54,6 +56,7 @@ import org.apache.cassandra.concurrent.DebuggableTask.RunnableDebuggableTask;
import org.apache.cassandra.concurrent.Stage;
import org.apache.cassandra.config.CassandraRelevantProperties;
import org.apache.cassandra.config.Config;
+import org.apache.cassandra.config.Config.NonSerialWriteStrategy;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.ConsistencyLevel;
@@ -90,8 +93,8 @@ import org.apache.cassandra.exceptions.QueryCancelledException;
import org.apache.cassandra.exceptions.ReadAbortException;
import org.apache.cassandra.exceptions.ReadFailureException;
import org.apache.cassandra.exceptions.ReadTimeoutException;
-import org.apache.cassandra.exceptions.RequestFailureException;
import org.apache.cassandra.exceptions.RequestFailure;
+import org.apache.cassandra.exceptions.RequestFailureException;
import org.apache.cassandra.exceptions.RequestTimeoutException;
import org.apache.cassandra.exceptions.UnavailableException;
import org.apache.cassandra.exceptions.WriteFailureException;
@@ -125,9 +128,18 @@ import org.apache.cassandra.schema.SchemaConstants;
import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.service.accord.AccordService;
+import org.apache.cassandra.service.accord.IAccordService;
+import org.apache.cassandra.service.accord.api.PartitionKey;
+import org.apache.cassandra.service.accord.txn.AccordUpdate;
+import org.apache.cassandra.service.accord.txn.TxnCondition;
import org.apache.cassandra.service.accord.txn.TxnData;
import org.apache.cassandra.service.accord.txn.TxnQuery;
import org.apache.cassandra.service.accord.txn.TxnRead;
+import org.apache.cassandra.service.accord.txn.TxnReferenceOperations;
+import org.apache.cassandra.service.accord.txn.TxnResult;
+import org.apache.cassandra.service.accord.txn.TxnUpdate;
+import org.apache.cassandra.service.accord.txn.TxnWrite;
+import org.apache.cassandra.service.consensus.migration.ConsensusRequestRouter;
import org.apache.cassandra.service.paxos.Ballot;
import org.apache.cassandra.service.paxos.Commit;
import org.apache.cassandra.service.paxos.ContentionStrategy;
@@ -137,6 +149,7 @@ import org.apache.cassandra.service.paxos.v1.PrepareCallback;
import org.apache.cassandra.service.paxos.v1.ProposeCallback;
import org.apache.cassandra.service.reads.AbstractReadExecutor;
import org.apache.cassandra.service.reads.ReadCallback;
+import org.apache.cassandra.service.reads.ReadCoordinator;
import org.apache.cassandra.service.reads.range.RangeCommands;
import org.apache.cassandra.service.reads.repair.ReadRepair;
import org.apache.cassandra.tcm.ClusterMetadata;
@@ -155,10 +168,10 @@ import org.apache.cassandra.utils.TimeUUID;
import org.apache.cassandra.utils.concurrent.CountDownLatch;
import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException;
+import static com.google.common.base.Preconditions.checkNotNull;
import static com.google.common.collect.Iterables.concat;
import static java.util.concurrent.TimeUnit.MILLISECONDS;
import static java.util.concurrent.TimeUnit.NANOSECONDS;
-import static org.apache.cassandra.config.Config.LegacyPaxosStrategy.accord;
import static org.apache.cassandra.db.ConsistencyLevel.SERIAL;
import static org.apache.cassandra.metrics.ClientRequestsMetricsHolder.casReadMetrics;
import static org.apache.cassandra.metrics.ClientRequestsMetricsHolder.casWriteMetrics;
@@ -177,6 +190,11 @@ import static org.apache.cassandra.net.Verb.PAXOS_PROPOSE_REQ;
import static org.apache.cassandra.net.Verb.SCHEMA_VERSION_REQ;
import static org.apache.cassandra.net.Verb.TRUNCATE_REQ;
import static org.apache.cassandra.service.BatchlogResponseHandler.BatchlogCleanup;
+import static org.apache.cassandra.service.StorageProxy.ConsensusAttemptResult.RETRY_NEW_PROTOCOL;
+import static org.apache.cassandra.service.StorageProxy.ConsensusAttemptResult.casResult;
+import static org.apache.cassandra.service.StorageProxy.ConsensusAttemptResult.serialReadResult;
+import static org.apache.cassandra.service.accord.txn.TxnResult.Kind.retry_new_protocol;
+import static org.apache.cassandra.service.consensus.migration.ConsensusRequestRouter.ConsensusRoutingDecision;
import static org.apache.cassandra.service.paxos.Ballot.Flag.GLOBAL;
import static org.apache.cassandra.service.paxos.Ballot.Flag.LOCAL;
import static org.apache.cassandra.service.paxos.BallotGenerator.Global.nextBallot;
@@ -322,6 +340,7 @@ public class StorageProxy implements StorageProxyMBean
Dispatcher.RequestTime requestTime)
throws UnavailableException, IsBootstrappingException, RequestFailureException, RequestTimeoutException, InvalidRequestException, CasWriteUnknownResultException
{
+ TableMetadata metadata = Schema.instance.validateTable(keyspaceName, cfName);
if (DatabaseDescriptor.getPartitionDenylistEnabled() && DatabaseDescriptor.getDenylistWritesEnabled() && !partitionDenylist.isKeyPermitted(keyspaceName, cfName, key.getKey()))
{
denylistMetrics.incrementWritesRejected();
@@ -329,34 +348,61 @@ public class StorageProxy implements StorageProxyMBean
key, keyspaceName, cfName));
}
- if (DatabaseDescriptor.getLegacyPaxosStrategy() == accord)
+ ConsensusAttemptResult lastAttemptResult;
+ do
{
- TxnData data = AccordService.instance().coordinate(request.toAccordTxn(clientState, nowInSeconds), consistencyForPaxos);
- return request.toCasResult(data);
- }
- else
- {
- return (Paxos.useV2() || keyspaceName.equals(SchemaConstants.METADATA_KEYSPACE_NAME))
- ? Paxos.cas(key, request, consistencyForPaxos, consistencyForCommit, clientState)
- : legacyCas(keyspaceName, cfName, key, request, consistencyForPaxos, consistencyForCommit, clientState, nowInSeconds, requestTime);
- }
+ ConsensusRoutingDecision decision = consensusRouting(metadata, key, consistencyForPaxos, requestTime, true);
+ switch (decision)
+ {
+ case paxosV2:
+ lastAttemptResult = Paxos.cas(key,
+ request,
+ consistencyForPaxos,
+ consistencyForCommit,
+ clientState,
+ requestTime);
+ break;
+ case paxosV1:
+ lastAttemptResult = legacyCas(metadata,
+ key,
+ request,
+ consistencyForPaxos,
+ consistencyForCommit,
+ clientState,
+ nowInSeconds,
+ requestTime);
+ break;
+ case accord:
+ Txn txn = request.toAccordTxn(consistencyForPaxos,
+ consistencyForCommit,
+ clientState,
+ nowInSeconds);
+ IAccordService accordService = AccordService.instance();
+ accordService.maybeConvertKeyspacesToAccord(txn);
+ TxnResult txnResult = accordService.coordinate(txn,
+ consistencyForPaxos,
+ requestTime);
+ lastAttemptResult = request.toCasResult(txnResult);
+ break;
+ default:
+ throw new IllegalStateException("Unsupported consensus " + decision);
+ }
+ } while (lastAttemptResult.shouldRetryOnNewConsensusProtocol);
+ return lastAttemptResult.casResult;
}
- public static RowIterator legacyCas(String keyspaceName,
- String cfName,
- DecoratedKey key,
- CASRequest request,
- ConsistencyLevel consistencyForPaxos,
- ConsistencyLevel consistencyForCommit,
- ClientState clientState,
- long nowInSeconds,
- Dispatcher.RequestTime requestTime)
+ private static ConsensusAttemptResult legacyCas(TableMetadata metadata,
+ DecoratedKey key,
+ CASRequest request,
+ ConsistencyLevel consistencyForPaxos,
+ ConsistencyLevel consistencyForCommit,
+ ClientState clientState,
+ long nowInSeconds,
+ Dispatcher.RequestTime requestTime)
throws UnavailableException, IsBootstrappingException, RequestFailureException, RequestTimeoutException, InvalidRequestException
{
try
{
- TableMetadata metadata = Schema.instance.validateTable(keyspaceName, cfName);
-
Function> updateProposer = ballot ->
{
// read the current values and check they validate the conditions
@@ -374,7 +420,7 @@ public class StorageProxy implements StorageProxyMBean
{
Tracing.trace("CAS precondition does not match current values {}", current);
casWriteMetrics.conditionNotMet.inc();
- return Pair.create(PartitionUpdate.emptyUpdate(metadata, key), current.rowIterator());
+ return Pair.create(PartitionUpdate.emptyUpdate(metadata, key), current.rowIterator(false));
}
// Create the desired updates
@@ -399,15 +445,14 @@ public class StorageProxy implements StorageProxyMBean
return Pair.create(updates, null);
};
- return doPaxos(metadata,
- key,
- consistencyForPaxos,
- consistencyForCommit,
- consistencyForCommit,
- requestTime,
- casWriteMetrics,
- updateProposer);
-
+ return casResult(doPaxos(metadata,
+ key,
+ consistencyForPaxos,
+ consistencyForCommit,
+ consistencyForCommit,
+ requestTime,
+ casWriteMetrics,
+ updateProposer));
}
catch (CasWriteUnknownResultException e)
{
@@ -1165,6 +1210,7 @@ public class StorageProxy implements StorageProxyMBean
Collection augmented = TriggerExecutor.instance.execute(mutations);
+ String keyspaceName = mutations.iterator().next().getKeyspaceName();
boolean updatesView = Keyspace.open(mutations.iterator().next().getKeyspaceName())
.viewManager
.updatesAffectView(mutations, true);
@@ -1172,8 +1218,10 @@ public class StorageProxy implements StorageProxyMBean
long size = IMutation.dataSize(mutations);
writeMetrics.mutationSize.update(size);
writeMetricsForLevel(consistencyLevel).mutationSize.update(size);
-
- if (augmented != null)
+ NonSerialWriteStrategy nonSerialWriteStrategy = DatabaseDescriptor.getNonSerialWriteStrategy();
+ if (nonSerialWriteStrategy.writesThroughAccord && !SchemaConstants.getSystemKeyspaces().contains(keyspaceName))
+ mutateWithAccord(augmented != null ? augmented : mutations, consistencyLevel, requestTime, nonSerialWriteStrategy);
+ else if (augmented != null)
mutateAtomically(augmented, consistencyLevel, updatesView, requestTime);
else
{
@@ -1184,6 +1232,29 @@ public class StorageProxy implements StorageProxyMBean
}
}
+ private static void mutateWithAccord(Collection extends IMutation> iMutations, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime, Config.NonSerialWriteStrategy nonSerialWriteStrategy)
+ {
+ int fragmentIndex = 0;
+ List fragments = new ArrayList<>(iMutations.size());
+ List partitionKeys = new ArrayList<>(iMutations.size());
+ for (IMutation mutation : iMutations)
+ {
+ for (PartitionUpdate update : mutation.getPartitionUpdates())
+ {
+ PartitionKey pk = PartitionKey.of(update);
+ partitionKeys.add(pk);
+ fragments.add(new TxnWrite.Fragment(PartitionKey.of(update), fragmentIndex++, update, TxnReferenceOperations.empty()));
+ }
+ }
+ // Potentially ignore commit consistency level if the strategy specifies accord and not migration
+ ConsistencyLevel clForCommit = nonSerialWriteStrategy.commitCLForStrategy(consistencyLevel);
+ AccordUpdate update = new TxnUpdate(fragments, TxnCondition.none(), clForCommit);
+ Txn.InMemory txn = new Txn.InMemory(Keys.of(partitionKeys), TxnRead.EMPTY, TxnQuery.EMPTY, update);
+ IAccordService accordService = AccordService.instance();
+ accordService.maybeConvertKeyspacesToAccord(txn);
+ accordService.coordinate(txn, consistencyLevel, requestTime);
+ }
+
/**
* See mutate. Adds additional steps before and after writing a batch.
* Before writing the batch (but after doing availability check against the FD for the row replicas):
@@ -1602,7 +1673,7 @@ public class StorageProxy implements StorageProxyMBean
if (insertLocal)
{
- Preconditions.checkNotNull(localReplica);
+ checkNotNull(localReplica);
performLocally(stage, localReplica, mutation::apply, responseHandler, mutation, requestTime);
}
@@ -1881,43 +1952,68 @@ public class StorageProxy implements StorageProxyMBean
return metadata.myNodeState() == NodeState.JOINED;
}
+ private static ConsensusRoutingDecision consensusRouting(TableMetadata metadata, DecoratedKey partitionKey, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime, boolean isForWrite)
+ {
+ if (metadata.keyspace.equals(SchemaConstants.METADATA_KEYSPACE_NAME))
+ return ConsensusRoutingDecision.paxosV2;
+ return ConsensusRequestRouter.instance.routeAndMaybeMigrate(partitionKey,
+ metadata.id,
+ consistencyLevel,
+ requestTime,
+ DatabaseDescriptor.getCasContentionTimeout(NANOSECONDS),
+ isForWrite);
+ }
+
private static PartitionIterator readWithConsensus(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime)
throws InvalidRequestException, UnavailableException, ReadFailureException, ReadTimeoutException
{
- // TCM explicitly relies on paxos and doesn't work with accord
- if (DatabaseDescriptor.getLegacyPaxosStrategy() == accord && !group.metadata().keyspace.equals(SchemaConstants.METADATA_KEYSPACE_NAME))
+ ConsensusAttemptResult lastResult;
+ do
{
- return readWithAccord(group, consistencyLevel);
- }
- else
- {
- return readWithPaxos(group, consistencyLevel, requestTime);
- }
+ SinglePartitionReadCommand command = group.queries.get(0);
+ ConsensusRoutingDecision decision = consensusRouting(group.metadata(), command.partitionKey(), consistencyLevel, requestTime, false);
+ switch (decision)
+ {
+ case paxosV2:
+ lastResult = Paxos.read(group, consistencyLevel, requestTime);
+ break;
+ case paxosV1:
+ lastResult = legacyReadWithPaxos(group, consistencyLevel, requestTime);
+ break;
+ case accord:
+ lastResult = readWithAccord(group, consistencyLevel, requestTime);
+ break;
+ default:
+ throw new IllegalStateException("Unsupported consensus " + decision);
+ }
+ } while (lastResult.shouldRetryOnNewConsensusProtocol);
+ return lastResult.serialReadResult;
}
- private static PartitionIterator readWithPaxos(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime)
- throws InvalidRequestException, UnavailableException, ReadFailureException, ReadTimeoutException
- {
- return (Paxos.useV2() || group.metadata().keyspace.equals(SchemaConstants.METADATA_KEYSPACE_NAME))
- ? Paxos.read(group, consistencyLevel, requestTime)
- : legacyReadWithPaxos(group, consistencyLevel, requestTime);
- }
-
- private static PartitionIterator readWithAccord(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel)
+ private static ConsensusAttemptResult readWithAccord(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime)
{
if (group.queries.size() > 1)
throw new InvalidRequestException("SERIAL/LOCAL_SERIAL consistency may only be requested for one partition at a time");
- TxnRead read = TxnRead.createSerialRead(group.queries.get(0));
+ SinglePartitionReadCommand readCommand = group.queries.get(0);
+ // If the non-SERIAL write strategy is sending all writes through Accord there is no need to use the supplied consistency
+ // level since Accord will manage reading safely
+ consistencyLevel = DatabaseDescriptor.getNonSerialWriteStrategy().readCLForStrategy(consistencyLevel);
+ TxnRead read = TxnRead.createSerialRead(readCommand, consistencyLevel);
Txn txn = new Txn.InMemory(read.keys(), read, TxnQuery.ALL);
- TxnData data = AccordService.instance().coordinate(txn, consistencyLevel);
+ IAccordService accordService = AccordService.instance();
+ accordService.maybeConvertKeyspacesToAccord(txn);
+ TxnResult txnResult = accordService.coordinate(txn, consistencyLevel, requestTime);
+ if (txnResult.kind() == retry_new_protocol)
+ return RETRY_NEW_PROTOCOL;
+ TxnData data = (TxnData)txnResult;
FilteredPartition partition = data.get(TxnRead.SERIAL_READ);
if (partition != null)
- return PartitionIterators.singletonIterator(partition.rowIterator());
+ return serialReadResult(PartitionIterators.singletonIterator(partition.rowIterator(readCommand.isReversed())));
else
- return EmptyIterators.partition();
+ return serialReadResult(EmptyIterators.partition());
}
- private static PartitionIterator legacyReadWithPaxos(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime)
+ private static ConsensusAttemptResult legacyReadWithPaxos(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime)
throws InvalidRequestException, UnavailableException, ReadFailureException, ReadTimeoutException
{
long start = nanoTime();
@@ -1930,7 +2026,6 @@ public class StorageProxy implements StorageProxyMBean
// calculate the blockFor before repair any paxos round to avoid RS being altered in between.
int blockForRead = consistencyLevel.blockFor(Keyspace.open(metadata.keyspace).getReplicationStrategy());
- PartitionIterator result = null;
try
{
final ConsistencyLevel consistencyForReplayCommitsOrFetch = consistencyLevel == ConsistencyLevel.LOCAL_SERIAL
@@ -1967,7 +2062,7 @@ public class StorageProxy implements StorageProxyMBean
throw new ReadFailureException(consistencyLevel, e.received, e.blockFor, false, e.failureReasonByEndpoint);
}
- result = fetchRows(group.queries, consistencyForReplayCommitsOrFetch, requestTime);
+ return serialReadResult(fetchRows(group.queries, consistencyForReplayCommitsOrFetch, ReadCoordinator.DEFAULT, requestTime));
}
catch (UnavailableException e)
{
@@ -2011,18 +2106,16 @@ public class StorageProxy implements StorageProxyMBean
readMetricsForLevel(consistencyLevel).addNano(latency);
Keyspace.open(metadata.keyspace).getColumnFamilyStore(metadata.name).metric.coordinatorReadLatency.update(latency, TimeUnit.NANOSECONDS);
}
-
- return result;
}
@SuppressWarnings("resource")
- private static PartitionIterator readRegular(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime)
+ public static PartitionIterator readRegular(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel, ReadCoordinator coordinator, Dispatcher.RequestTime requestTime)
throws UnavailableException, ReadFailureException, ReadTimeoutException
{
long start = nanoTime();
try
{
- PartitionIterator result = fetchRows(group.queries, consistencyLevel, requestTime);
+ PartitionIterator result = fetchRows(group.queries, consistencyLevel, coordinator, requestTime);
// Note that the only difference between the command in a group must be the partition key on which
// they applied.
boolean enforceStrictLiveness = group.queries.get(0).metadata().enforceStrictLiveness();
@@ -2072,6 +2165,11 @@ public class StorageProxy implements StorageProxyMBean
}
}
+ public static PartitionIterator readRegular(SinglePartitionReadCommand.Group group, ConsistencyLevel consistencyLevel, Dispatcher.RequestTime requestTime)
+ {
+ return readRegular(group, consistencyLevel, ReadCoordinator.DEFAULT, requestTime);
+ }
+
public static void recordReadRegularAbort(ConsistencyLevel consistencyLevel, Throwable cause)
{
readMetrics.markAbort(cause);
@@ -2119,6 +2217,7 @@ public class StorageProxy implements StorageProxyMBean
*/
private static PartitionIterator fetchRows(List commands,
ConsistencyLevel consistencyLevel,
+ ReadCoordinator coordinator,
Dispatcher.RequestTime requestTime)
throws UnavailableException, ReadFailureException, ReadTimeoutException
{
@@ -2131,7 +2230,7 @@ public class StorageProxy implements StorageProxyMBean
// for type of speculation we'll use in this read
for (int i=0; i anyOutOfRangeOpsRecorded
= keyspace -> keyspace.metric.outOfRangeTokenReads.getCount() > 0
|| keyspace.metric.outOfRangeTokenWrites.getCount() > 0
@@ -1664,6 +1676,49 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
}
}
+ @Override
+ public void migrateConsensusProtocol(@Nonnull String targetProtocol,
+ @Nonnull List keyspaceNames,
+ @Nullable List maybeTableNames,
+ @Nullable String maybeRangesStr)
+ {
+ checkNotNull(targetProtocol, "targetProtocol is null");
+ checkArgument(!keyspaceNames.contains(SchemaConstants.METADATA_KEYSPACE_NAME));
+ startMigrationToConsensusProtocol(targetProtocol, keyspaceNames, Optional.ofNullable(maybeTableNames), Optional.ofNullable(maybeRangesStr));
+ }
+
+ @Override
+ public List finishConsensusMigration(@Nonnull String keyspace,
+ @Nullable List maybeTableNames,
+ @Nullable String maybeRangesStr)
+ {
+ checkArgument(!keyspace.equals(SchemaConstants.METADATA_KEYSPACE_NAME));
+ return finishMigrationToConsensusProtocol(keyspace, Optional.ofNullable(maybeTableNames), Optional.ofNullable(maybeRangesStr));
+ }
+
+ @Override
+ public void setConsensusMigrationTargetProtocol(@Nonnull String targetProtocol,
+ @Nullable List keyspaceNames,
+ @Nullable List maybeTableNames)
+ {
+ checkNotNull(targetProtocol, "targetProtocol is null");
+ checkNotNull(keyspaceNames, "keyspaceNames is null");
+ checkArgument(!keyspaceNames.contains(SchemaConstants.METADATA_KEYSPACE_NAME));
+
+ ConsensusTableMigrationState.setConsensusMigrationTargetProtocol(targetProtocol, keyspaceNames, Optional.ofNullable(maybeTableNames));
+ }
+
+ @Override
+ public String listConsensusMigrations(@Nullable Set keyspaceNames,
+ @Nullable Set tableNames,
+ @Nonnull String format)
+ {
+ ClusterMetadata cm = ClusterMetadata.current();
+ ConsensusMigrationState snapshot = cm.consensusMigrationState;
+ Map snapshotAsMap = snapshot.toMap(keyspaceNames, tableNames);
+ return pojoMapToString(snapshotAsMap, format);
+ }
+
public Map> getConcurrency(List stageNames)
{
Stream stageStream = stageNames.isEmpty() ? stream(Stage.values()) : stageNames.stream().map(Stage::fromPoolName);
@@ -3978,11 +4033,21 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
// Never ever do this at home. Used by tests.
@VisibleForTesting
- public IPartitioner setPartitionerUnsafe(IPartitioner newPartitioner)
+ public void setPartitionerUnsafe(IPartitioner newPartitioner)
{
- IPartitioner oldPartitioner = DatabaseDescriptor.setPartitionerUnsafe(newPartitioner);
+ checkNotNull(newPartitioner, "newPartitioner is null");
+ checkState(originalPartitioner == null, "Already changed the partitioner without resetting");
+ originalPartitioner = DatabaseDescriptor.setPartitionerUnsafe(newPartitioner);
valueFactory = new VersionedValue.VersionedValueFactory(newPartitioner);
- return oldPartitioner;
+ }
+
+ @VisibleForTesting
+ public void resetPartitionerUnsafe()
+ {
+ checkState(originalPartitioner != null, "Original partitioner was never changed");
+ DatabaseDescriptor.setPartitionerUnsafe(originalPartitioner);
+ valueFactory = new VersionedValueFactory(originalPartitioner);
+ originalPartitioner = null;
}
public void truncate(String keyspace, String table) throws TimeoutException, IOException
@@ -4165,6 +4230,16 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
return Lists.newArrayList(Schema.instance.distributedKeyspaces().names());
}
+ @Override
+ public List getAccordManagedKeyspaces()
+ {
+ // TODO (review) These are really just the ones Accord is aware of not necessarily managed
+ Set keyspaces = Schema.instance.getNonLocalStrategyKeyspaces().names();
+ return keyspaces.stream()
+ .filter(AccordService.instance()::isAccordManagedKeyspace)
+ .collect(toList());
+ }
+
public Map getViewBuildStatuses(String keyspace, String view, boolean withPort)
{
Map coreViewStatus = SystemDistributedKeyspace.viewStatus(keyspace, view);
@@ -4994,7 +5069,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
archiveCommand = archiveCommand != null ? archiveCommand : fqlOptions.archive_command;
maxArchiveRetries = maxArchiveRetries != Integer.MIN_VALUE ? maxArchiveRetries : fqlOptions.max_archive_retries;
- Preconditions.checkNotNull(path, "cassandra.yaml did not set log_dir and not set as parameter");
+ checkNotNull(path, "cassandra.yaml did not set log_dir and not set as parameter");
FullQueryLogger.instance.enableWithoutClean(File.getPath(path), rollCycle, blocking, maxQueueWeight, maxLogSize, archiveCommand, maxArchiveRetries);
}
@@ -5359,7 +5434,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
public void setRepairRpcTimeout(Long timeoutInMillis)
{
- Preconditions.checkState(timeoutInMillis > 0);
+ checkState(timeoutInMillis > 0);
DatabaseDescriptor.setRepairRpcTimeout(timeoutInMillis);
logger.info("RepairRpcTimeout set to {}ms via JMX", timeoutInMillis);
}
diff --git a/src/java/org/apache/cassandra/service/StorageServiceMBean.java b/src/java/org/apache/cassandra/service/StorageServiceMBean.java
index 6e8a2ad449..d11ca1bfb4 100644
--- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java
+++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java
@@ -27,6 +27,7 @@ import java.util.Map;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeoutException;
+import javax.annotation.Nonnull;
import javax.annotation.Nullable;
import javax.management.NotificationEmitter;
import javax.management.openmbean.CompositeData;
@@ -1141,6 +1142,23 @@ public interface StorageServiceMBean extends NotificationEmitter
public String getBootstrapState();
void abortBootstrap(String nodeId, String endpoint);
+ void migrateConsensusProtocol(@Nonnull String targetProtocol,
+ @Nullable List keyspaceNames,
+ @Nullable List maybeTableNames,
+ @Nullable String maybeRangesStr);
+
+ List finishConsensusMigration(@Nonnull String keyspace,
+ @Nullable List maybeTableNames,
+ @Nullable String maybeRangesStr);
+
+ void setConsensusMigrationTargetProtocol(@Nonnull String targetProtocol,
+ @Nullable List keyspaceNames,
+ @Nullable List maybeTableNames);
+
+ String listConsensusMigrations(@Nullable Set keyspaceNames, @Nullable Set tableNames, @Nonnull String format);
+
+ List getAccordManagedKeyspaces();
+
/** Gets the concurrency settings for processing stages*/
static class StageConcurrency implements Serializable
{
@@ -1188,6 +1206,7 @@ public interface StorageServiceMBean extends NotificationEmitter
/**
* Start the fully query logger.
+ *
* @param path Path where the full query log will be stored. If null cassandra.yaml value is used.
* @param rollCycle How often to create a new file for query data (MINUTELY, DAILY, HOURLY)
* @param blocking Whether threads submitting queries to the query log should block if they can't be drained to the filesystem or alternatively drops samples and log
diff --git a/src/java/org/apache/cassandra/service/accord/AccordCachingState.java b/src/java/org/apache/cassandra/service/accord/AccordCachingState.java
index 994e551f79..d7bce189d2 100644
--- a/src/java/org/apache/cassandra/service/accord/AccordCachingState.java
+++ b/src/java/org/apache/cassandra/service/accord/AccordCachingState.java
@@ -17,8 +17,6 @@
*/
package org.apache.cassandra.service.accord;
-import java.util.Collections;
-import java.util.Set;
import java.util.concurrent.Callable;
import java.util.function.BiFunction;
import java.util.function.Function;
@@ -27,7 +25,7 @@ import java.util.function.ToLongFunction;
import com.google.common.primitives.Ints;
import accord.local.Command.TransientListener;
-import accord.utils.DeterministicIdentitySet;
+import accord.local.Listeners;
import accord.utils.IntrusiveLinkedListNode;
import accord.utils.async.AsyncChain;
import accord.utils.async.AsyncResults.RunnableResult;
@@ -35,14 +33,14 @@ import org.apache.cassandra.concurrent.ExecutorPlus;
import org.apache.cassandra.utils.ObjectSizes;
import static java.lang.String.format;
-import static org.apache.cassandra.service.accord.AccordCachingState.Status.UNINITIALIZED;
-import static org.apache.cassandra.service.accord.AccordCachingState.Status.LOADING;
-import static org.apache.cassandra.service.accord.AccordCachingState.Status.LOADED;
+import static org.apache.cassandra.service.accord.AccordCachingState.Status.EVICTED;
import static org.apache.cassandra.service.accord.AccordCachingState.Status.FAILED_TO_LOAD;
+import static org.apache.cassandra.service.accord.AccordCachingState.Status.FAILED_TO_SAVE;
+import static org.apache.cassandra.service.accord.AccordCachingState.Status.LOADED;
+import static org.apache.cassandra.service.accord.AccordCachingState.Status.LOADING;
import static org.apache.cassandra.service.accord.AccordCachingState.Status.MODIFIED;
import static org.apache.cassandra.service.accord.AccordCachingState.Status.SAVING;
-import static org.apache.cassandra.service.accord.AccordCachingState.Status.FAILED_TO_SAVE;
-import static org.apache.cassandra.service.accord.AccordCachingState.Status.EVICTED;
+import static org.apache.cassandra.service.accord.AccordCachingState.Status.UNINITIALIZED;
/**
* Global (per CommandStore) state of a cached entity (Command or CommandsForKey).
@@ -61,7 +59,7 @@ public class AccordCachingState