mirror of https://github.com/apache/cassandra
CEP-15 (C*) When a host replacement happens don't loose the peer mapping right away (#3575)
patch by David Capwell; reviewed by Blake Eggleston for CASSANDRA-18764
This commit is contained in:
parent
a6bb08d926
commit
76994f8196
|
|
@ -1 +1 @@
|
|||
Subproject commit 7c15f3a6203939bc6cb398e538df1ca3557cbe03
|
||||
Subproject commit 2ad55e03c43ce074cdf5e36cfa14cb4278c2dc0f
|
||||
|
|
@ -272,36 +272,36 @@ public enum Verb
|
|||
// accord
|
||||
ACCORD_SIMPLE_RSP (119, P2, writeTimeout, REQUEST_RESPONSE, () -> EnumSerializer.simpleReply, RESPONSE_HANDLER ),
|
||||
ACCORD_PRE_ACCEPT_RSP (121, P2, writeTimeout, REQUEST_RESPONSE, () -> PreacceptSerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_PRE_ACCEPT_REQ (120, P2, writeTimeout, IMMEDIATE, () -> PreacceptSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_PRE_ACCEPT_RSP ),
|
||||
ACCORD_PRE_ACCEPT_REQ (120, P2, writeTimeout, IMMEDIATE, () -> PreacceptSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_PRE_ACCEPT_RSP ),
|
||||
ACCORD_ACCEPT_RSP (124, P2, writeTimeout, REQUEST_RESPONSE, () -> AcceptSerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_ACCEPT_REQ (122, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_ACCEPT_RSP ),
|
||||
ACCORD_ACCEPT_INVALIDATE_REQ (123, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.invalidate, () -> AccordService.instance().verbHandler(), ACCORD_ACCEPT_RSP ),
|
||||
ACCORD_ACCEPT_REQ (122, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_ACCEPT_RSP ),
|
||||
ACCORD_ACCEPT_INVALIDATE_REQ (123, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.invalidate, AccordService::verbHandlerOrNoop, ACCORD_ACCEPT_RSP ),
|
||||
ACCORD_READ_RSP (126, P2, writeTimeout, REQUEST_RESPONSE, () -> ReadDataSerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_READ_REQ (125, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_READ_RSP ),
|
||||
ACCORD_COMMIT_REQ (127, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_READ_RSP ),
|
||||
ACCORD_COMMIT_INVALIDATE_REQ (128, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.invalidate, () -> AccordService.instance().verbHandler() ),
|
||||
ACCORD_READ_REQ (125, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_READ_RSP ),
|
||||
ACCORD_COMMIT_REQ (127, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_READ_RSP ),
|
||||
ACCORD_COMMIT_INVALIDATE_REQ (128, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.invalidate, AccordService::verbHandlerOrNoop ),
|
||||
ACCORD_APPLY_RSP (130, P2, writeTimeout, REQUEST_RESPONSE, () -> ApplySerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_APPLY_REQ (129, P2, writeTimeout, IMMEDIATE, () -> ApplySerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_APPLY_RSP ),
|
||||
ACCORD_APPLY_REQ (129, P2, writeTimeout, IMMEDIATE, () -> ApplySerializers.request, AccordService::verbHandlerOrNoop, ACCORD_APPLY_RSP ),
|
||||
ACCORD_BEGIN_RECOVER_RSP (132, P2, writeTimeout, REQUEST_RESPONSE, () -> RecoverySerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_BEGIN_RECOVER_REQ (131, P2, writeTimeout, IMMEDIATE, () -> RecoverySerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_BEGIN_RECOVER_RSP ),
|
||||
ACCORD_BEGIN_RECOVER_REQ (131, P2, writeTimeout, IMMEDIATE, () -> RecoverySerializers.request, AccordService::verbHandlerOrNoop, ACCORD_BEGIN_RECOVER_RSP ),
|
||||
ACCORD_BEGIN_INVALIDATE_RSP (134, P2, writeTimeout, REQUEST_RESPONSE, () -> BeginInvalidationSerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_BEGIN_INVALIDATE_REQ (133, P2, writeTimeout, IMMEDIATE, () -> BeginInvalidationSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_BEGIN_INVALIDATE_RSP ),
|
||||
ACCORD_BEGIN_INVALIDATE_REQ (133, P2, writeTimeout, IMMEDIATE, () -> BeginInvalidationSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_BEGIN_INVALIDATE_RSP ),
|
||||
ACCORD_WAIT_ON_COMMIT_RSP (136, P2, writeTimeout, REQUEST_RESPONSE, () -> WaitOnCommitSerializer.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_WAIT_ON_COMMIT_REQ (135, P2, writeTimeout, IMMEDIATE, () -> WaitOnCommitSerializer.request, () -> AccordService.instance().verbHandler(), ACCORD_WAIT_ON_COMMIT_RSP ),
|
||||
ACCORD_WAIT_ON_APPLY_REQ (137, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.waitOnApply, () -> AccordService.instance().verbHandler(), ACCORD_READ_RSP ),
|
||||
ACCORD_INFORM_OF_TXN_REQ (138, P2, writeTimeout, IMMEDIATE, () -> InformOfTxnIdSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_INFORM_HOME_DURABLE_REQ (139, P2, writeTimeout, IMMEDIATE, () -> InformHomeDurableSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_INFORM_DURABLE_REQ (140, P2, writeTimeout, IMMEDIATE, () -> InformDurableSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_WAIT_ON_COMMIT_REQ (135, P2, writeTimeout, IMMEDIATE, () -> WaitOnCommitSerializer.request, AccordService::verbHandlerOrNoop, ACCORD_WAIT_ON_COMMIT_RSP ),
|
||||
ACCORD_WAIT_ON_APPLY_REQ (137, P2, writeTimeout, IMMEDIATE, () -> ReadDataSerializers.waitOnApply, AccordService::verbHandlerOrNoop, ACCORD_READ_RSP ),
|
||||
ACCORD_INFORM_OF_TXN_REQ (138, P2, writeTimeout, IMMEDIATE, () -> InformOfTxnIdSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_INFORM_HOME_DURABLE_REQ (139, P2, writeTimeout, IMMEDIATE, () -> InformHomeDurableSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_INFORM_DURABLE_REQ (140, P2, writeTimeout, IMMEDIATE, () -> InformDurableSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_CHECK_STATUS_RSP (142, P2, writeTimeout, REQUEST_RESPONSE, () -> CheckStatusSerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_CHECK_STATUS_REQ (141, P2, writeTimeout, IMMEDIATE, () -> CheckStatusSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_CHECK_STATUS_RSP ),
|
||||
ACCORD_CHECK_STATUS_REQ (141, P2, writeTimeout, IMMEDIATE, () -> CheckStatusSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_CHECK_STATUS_RSP ),
|
||||
ACCORD_GET_DEPS_RSP (144, P2, writeTimeout, REQUEST_RESPONSE, () -> GetDepsSerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_GET_DEPS_REQ (143, P2, writeTimeout, IMMEDIATE, () -> GetDepsSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_GET_DEPS_RSP ),
|
||||
ACCORD_GET_DEPS_REQ (143, P2, writeTimeout, IMMEDIATE, () -> GetDepsSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_GET_DEPS_RSP ),
|
||||
ACCORD_FETCH_DATA_RSP (146, P2, repairTimeout,REQUEST_RESPONSE, () -> FetchSerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_FETCH_DATA_REQ (145, P2, repairTimeout,IMMEDIATE, () -> FetchSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_FETCH_DATA_RSP ),
|
||||
ACCORD_SET_SHARD_DURABLE_REQ (147, P2, writeTimeout, IMMEDIATE, () -> SetDurableSerializers.shardDurable, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_SET_GLOBALLY_DURABLE_REQ (148, P2, writeTimeout, IMMEDIATE, () -> SetDurableSerializers.globallyDurable,() -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_FETCH_DATA_REQ (145, P2, repairTimeout,IMMEDIATE, () -> FetchSerializers.request, AccordService::verbHandlerOrNoop, ACCORD_FETCH_DATA_RSP ),
|
||||
ACCORD_SET_SHARD_DURABLE_REQ (147, P2, writeTimeout, IMMEDIATE, () -> SetDurableSerializers.shardDurable, AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_SET_GLOBALLY_DURABLE_REQ (148, P2, writeTimeout, IMMEDIATE, () -> SetDurableSerializers.globallyDurable,AccordService::verbHandlerOrNoop, ACCORD_SIMPLE_RSP ),
|
||||
ACCORD_QUERY_DURABLE_BEFORE_RSP (150, P2, writeTimeout, REQUEST_RESPONSE, () -> QueryDurableBeforeSerializers.reply, RESPONSE_HANDLER ),
|
||||
ACCORD_QUERY_DURABLE_BEFORE_REQ (149, P2, writeTimeout, IMMEDIATE, () -> QueryDurableBeforeSerializers.request,() -> AccordService.instance().verbHandler(), ACCORD_QUERY_DURABLE_BEFORE_RSP),
|
||||
ACCORD_QUERY_DURABLE_BEFORE_REQ (149, P2, writeTimeout, IMMEDIATE, () -> QueryDurableBeforeSerializers.request,AccordService::verbHandlerOrNoop, ACCORD_QUERY_DURABLE_BEFORE_RSP),
|
||||
|
||||
ACCORD_SYNC_NOTIFY_REQ (151, P2, writeTimeout, IMMEDIATE, () -> Notification.listSerializer, () -> AccordSyncPropagator.verbHandler, ACCORD_SIMPLE_RSP ),
|
||||
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ import java.util.stream.Collectors;
|
|||
import javax.annotation.Nullable;
|
||||
|
||||
import com.google.common.annotations.VisibleForTesting;
|
||||
import com.google.common.collect.Sets;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
@ -132,13 +133,17 @@ public class AccordConfigurationService extends AbstractConfigurationService<Acc
|
|||
Invariants.checkState(state == State.INITIALIZED, "Expected state to be INITIALIZED but was %s", state);
|
||||
state = State.LOADING;
|
||||
updateMapping(ClusterMetadata.current());
|
||||
EndpointMapping snapshot = mapping;
|
||||
diskState = AccordKeyspace.loadTopologies(((epoch, topology, syncStatus, pendingSyncNotify, remoteSyncComplete, closed, redundant) -> {
|
||||
if (topology != null)
|
||||
reportTopology(topology, syncStatus == SyncStatus.NOT_STARTED);
|
||||
|
||||
getOrCreateEpochState(epoch).setSyncStatus(syncStatus);
|
||||
if (syncStatus == SyncStatus.NOTIFYING)
|
||||
syncPropagator.reportSyncComplete(epoch, pendingSyncNotify, localId);
|
||||
{
|
||||
// TODO (expected, correctness): since this is loading old topologies, might see nodes no longer present (host replacement, decom, shrink, etc.); attempt to remove unknown nodes
|
||||
syncPropagator.reportSyncComplete(epoch, Sets.filter(pendingSyncNotify, snapshot::containsId), localId);
|
||||
}
|
||||
|
||||
remoteSyncComplete.forEach(id -> receiveRemoteSyncComplete(id, epoch));
|
||||
// TODO (now): disk doesn't get updated until we see our own notification, so there is an edge case where this instance notified others and fails in the middle, but Apply was already sent! This could leave partial closed/redudant accross the cluster
|
||||
|
|
@ -157,7 +162,7 @@ public class AccordConfigurationService extends AbstractConfigurationService<Acc
|
|||
@Override
|
||||
public InetAddressAndPort mappedEndpoint(Node.Id id)
|
||||
{
|
||||
return Invariants.nonNull(mapping.mappedEndpoint(id));
|
||||
return Invariants.nonNull(mapping.mappedEndpoint(id), "Unable to map node id %s to a InetAddressAndPort", id);
|
||||
}
|
||||
|
||||
@VisibleForTesting
|
||||
|
|
@ -175,7 +180,7 @@ public class AccordConfigurationService extends AbstractConfigurationService<Acc
|
|||
|
||||
synchronized void updateMapping(ClusterMetadata metadata)
|
||||
{
|
||||
updateMapping(AccordTopologyUtils.directoryToMapping(metadata.epoch.getEpoch(), metadata.directory));
|
||||
updateMapping(AccordTopologyUtils.directoryToMapping(mapping, metadata.epoch.getEpoch(), metadata.directory));
|
||||
}
|
||||
|
||||
private void reportMetadata(ClusterMetadata metadata)
|
||||
|
|
@ -220,11 +225,11 @@ public class AccordConfigurationService extends AbstractConfigurationService<Acc
|
|||
}
|
||||
|
||||
@Override
|
||||
protected synchronized void localSyncComplete(Topology topology)
|
||||
protected synchronized void localSyncComplete(Topology topology, boolean startSync)
|
||||
{
|
||||
long epoch = topology.epoch();
|
||||
EpochState epochState = getOrCreateEpochState(epoch);
|
||||
if (epochState.syncStatus != SyncStatus.NOT_STARTED)
|
||||
if (!startSync ||epochState.syncStatus != SyncStatus.NOT_STARTED)
|
||||
return;
|
||||
|
||||
Set<Node.Id> notify = topology.nodes().stream().filter(i -> !localId.equals(i)).collect(Collectors.toSet());
|
||||
|
|
|
|||
|
|
@ -163,12 +163,24 @@ public class AccordService implements IAccordService, Shutdownable
|
|||
}
|
||||
};
|
||||
|
||||
private static Node.Id localId = null;
|
||||
private static volatile Node.Id localId = null;
|
||||
|
||||
private static class Handle
|
||||
{
|
||||
public static final AccordService instance = new AccordService();
|
||||
}
|
||||
|
||||
public static boolean isSetup()
|
||||
{
|
||||
return localId != null;
|
||||
}
|
||||
|
||||
public static IVerbHandler<? extends Request> verbHandlerOrNoop()
|
||||
{
|
||||
if (!isSetup()) return ignore -> {};
|
||||
return instance().verbHandler();
|
||||
}
|
||||
|
||||
public static void startup(NodeId tcmId)
|
||||
{
|
||||
localId = AccordTopologyUtils.tcmIdToAccord(tcmId);
|
||||
|
|
|
|||
|
|
@ -59,7 +59,11 @@ import static org.apache.cassandra.utils.CollectionSerializers.newListSerializer
|
|||
*/
|
||||
public class AccordSyncPropagator
|
||||
{
|
||||
public static final IVerbHandler<List<Notification>> verbHandler = message -> AccordService.instance().receive(message);
|
||||
public static final IVerbHandler<List<Notification>> verbHandler = message -> {
|
||||
if (!AccordService.isSetup())
|
||||
return;
|
||||
AccordService.instance().receive(message);
|
||||
};
|
||||
|
||||
interface Listener
|
||||
{
|
||||
|
|
|
|||
|
|
@ -131,11 +131,16 @@ public class AccordTopologyUtils
|
|||
return new Topology(epoch.getEpoch(), shards.toArray(new Shard[0]));
|
||||
}
|
||||
|
||||
public static EndpointMapping directoryToMapping(long epoch, Directory directory)
|
||||
public static EndpointMapping directoryToMapping(EndpointMapping mapping, long epoch, Directory directory)
|
||||
{
|
||||
EndpointMapping.Builder builder = EndpointMapping.builder(epoch);
|
||||
for (NodeId id : directory.peerIds())
|
||||
builder.add(directory.endpoint(id), tcmIdToAccord(id));
|
||||
|
||||
// There are cases where nodes are removed from the cluster (host replacement, decom, etc.), but inflight events may still be happening;
|
||||
// keep the ids around so pending events do not fail with a mapping error
|
||||
for (Node.Id id : mapping.differenceIds(builder))
|
||||
builder.add(mapping.mappedEndpoint(id), id);
|
||||
return builder.build();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -18,9 +18,12 @@
|
|||
|
||||
package org.apache.cassandra.service.accord;
|
||||
|
||||
import java.util.Set;
|
||||
|
||||
import com.google.common.collect.BiMap;
|
||||
import com.google.common.collect.HashBiMap;
|
||||
import com.google.common.collect.ImmutableBiMap;
|
||||
import com.google.common.collect.Sets;
|
||||
|
||||
import accord.local.Node;
|
||||
import accord.utils.Invariants;
|
||||
|
|
@ -44,6 +47,16 @@ class EndpointMapping implements AccordEndpointMapper
|
|||
return epoch;
|
||||
}
|
||||
|
||||
public boolean containsId(Node.Id id)
|
||||
{
|
||||
return mapping.containsKey(id);
|
||||
}
|
||||
|
||||
public Set<Node.Id> differenceIds(Builder builder)
|
||||
{
|
||||
return Sets.difference(mapping.keySet(), builder.mapping.keySet());
|
||||
}
|
||||
|
||||
@Override
|
||||
public Node.Id mappedId(InetAddressAndPort endpoint)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -981,7 +981,10 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance
|
|||
() -> SharedExecutorPool.SHARED.shutdownAndWait(1L, MINUTES)
|
||||
);
|
||||
|
||||
error = parallelRun(error, executor, () -> AccordService.instance().shutdownAndWait(1l, MINUTES));
|
||||
error = parallelRun(error, executor, () -> {
|
||||
if (!AccordService.isSetup()) return;
|
||||
AccordService.instance().shutdownAndWait(1l, MINUTES);
|
||||
});
|
||||
|
||||
// CommitLog must shut down after Stage, or threads from the latter may attempt to use the former.
|
||||
// (ex. A Mutation stage thread may attempt to add a mutation to the CommitLog.)
|
||||
|
|
|
|||
|
|
@ -207,19 +207,19 @@ public class AccordConfigurationServiceTest
|
|||
Topology topology1 = new Topology(1, new Shard(AccordTopologyUtils.fullRange("ks"), ID_LIST, ID_SET));
|
||||
service.updateMapping(mappingForEpoch(ClusterMetadata.current().epoch.getEpoch() + 1));
|
||||
service.reportTopology(topology1);
|
||||
service.acknowledgeEpoch(EpochReady.done(1));
|
||||
service.acknowledgeEpoch(EpochReady.done(1), true);
|
||||
service.receiveRemoteSyncComplete(ID1, 1);
|
||||
service.receiveRemoteSyncComplete(ID2, 1);
|
||||
service.receiveRemoteSyncComplete(ID3, 1);
|
||||
|
||||
Topology topology2 = new Topology(2, new Shard(AccordTopologyUtils.fullRange("ks"), ID_LIST, of(ID1, ID2)));
|
||||
service.reportTopology(topology2);
|
||||
service.acknowledgeEpoch(EpochReady.done(2));
|
||||
service.acknowledgeEpoch(EpochReady.done(2), true);
|
||||
service.receiveRemoteSyncComplete(ID1, 2);
|
||||
|
||||
Topology topology3 = new Topology(3, new Shard(AccordTopologyUtils.fullRange("ks"), ID_LIST, of(ID1, ID2)));
|
||||
service.reportTopology(topology3);
|
||||
service.acknowledgeEpoch(EpochReady.done(3));
|
||||
service.acknowledgeEpoch(EpochReady.done(3), true);
|
||||
|
||||
AccordConfigurationService loaded = new AccordConfigurationService(ID1, new Messaging(), new MockFailureDetector());
|
||||
loaded.updateMapping(mappingForEpoch(ClusterMetadata.current().epoch.getEpoch() + 1));
|
||||
|
|
|
|||
|
|
@ -427,7 +427,7 @@ public class AccordSyncPropagatorTest
|
|||
}
|
||||
|
||||
@Override
|
||||
protected void localSyncComplete(Topology topology)
|
||||
protected void localSyncComplete(Topology topology, boolean startSync)
|
||||
{
|
||||
Set<Node.Id> notify = topology.nodes().stream().filter(i -> !localId.equals(i)).collect(Collectors.toSet());
|
||||
instances.get(localId).propagator.reportSyncComplete(topology.epoch(), notify, localId);
|
||||
|
|
|
|||
Loading…
Reference in New Issue