mirror of https://github.com/apache/cassandra
[CASSANDRA-20736] Make DataPlacements private on ClusterMetadata
This commit is contained in:
parent
bed10b544f
commit
360facb24c
|
|
@ -96,7 +96,7 @@ public final class CreateKeyspaceStatement extends AlterSchemaStatement
|
|||
// as we have as keys in metadata.placements to have a fast map lookup
|
||||
// ReplicationParams are immutable, so it is a safe optimization
|
||||
KeyspaceParams keyspaceParams = attrs.asNewKeyspaceParams();
|
||||
ReplicationParams replicationParams = metadata.placements.deduplicateReplicationParams(keyspaceParams.replication);
|
||||
ReplicationParams replicationParams = metadata.placements().deduplicateReplicationParams(keyspaceParams.replication);
|
||||
keyspaceParams = keyspaceParams.withSwapped(replicationParams);
|
||||
KeyspaceMetadata keyspaceMetadata = KeyspaceMetadata.create(keyspaceName, keyspaceParams);
|
||||
|
||||
|
|
|
|||
|
|
@ -195,6 +195,6 @@ public abstract class AbstractMutationVerbHandler<T extends IMutation> implement
|
|||
|
||||
private static VersionedEndpoints.ForToken writePlacements(ClusterMetadata metadata, String keyspace, DecoratedKey key)
|
||||
{
|
||||
return metadata.placements.get(metadata.schema.getKeyspace(keyspace).getMetadata().params.replication).writes.forToken(key.getToken());
|
||||
return metadata.placement(metadata.schema.getKeyspace(keyspace).getMetadata().params.replication).writes.forToken(key.getToken());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -240,8 +240,8 @@ public class ReadCommandVerbHandler implements IVerbHandler<ReadCommand>
|
|||
|
||||
private static Replica getLocalReplica(ClusterMetadata metadata, Token token, String keyspace)
|
||||
{
|
||||
return metadata.placements
|
||||
.get(metadata.schema.getKeyspaces().getNullable(keyspace).params.replication)
|
||||
return metadata
|
||||
.placement(metadata.schema.getKeyspaces().getNullable(keyspace).params.replication)
|
||||
.reads
|
||||
.forToken(token)
|
||||
.get()
|
||||
|
|
|
|||
|
|
@ -780,7 +780,7 @@ public class CompactionManager implements CompactionManagerMBean, ICompactionMan
|
|||
// we only consider write placements during cleanup as range movements always ensure
|
||||
// overlap between new replicas accepting reads and old replicas accepting writes
|
||||
ClusterMetadata cm = ClusterMetadata.current();
|
||||
DataPlacement placement = cm.placements.get(keyspace.getMetadata().params.replication);
|
||||
DataPlacement placement = cm.placement(keyspace.getMetadata().params.replication);
|
||||
InetAddressAndPort local = FBUtilities.getBroadcastAddressAndPort();
|
||||
RangesAtEndpoint localWrites = placement.writes.byEndpoint().get(local);
|
||||
// TODO review: Hack to get local partitioner not to fail out because it's handled very poorly with data placements
|
||||
|
|
|
|||
|
|
@ -30,6 +30,7 @@ import org.apache.cassandra.locator.Replica;
|
|||
import org.apache.cassandra.schema.KeyspaceMetadata;
|
||||
import org.apache.cassandra.tcm.ClusterMetadata;
|
||||
import org.apache.cassandra.tcm.membership.Location;
|
||||
import org.apache.cassandra.tcm.ownership.DataPlacement;
|
||||
|
||||
public final class ViewUtils
|
||||
{
|
||||
|
|
@ -64,8 +65,9 @@ public final class ViewUtils
|
|||
Location local = metadata.locator.local();
|
||||
KeyspaceMetadata keyspaceMetadata = metadata.schema.getKeyspaces().getNullable(keyspace);
|
||||
|
||||
EndpointsForToken naturalBaseReplicas = metadata.placements.get(keyspaceMetadata.params.replication).reads.forToken(baseToken).get();
|
||||
EndpointsForToken naturalViewReplicas = metadata.placements.get(keyspaceMetadata.params.replication).reads.forToken(viewToken).get();
|
||||
DataPlacement placement = metadata.placement(keyspaceMetadata.params.replication);
|
||||
EndpointsForToken naturalBaseReplicas = placement.reads.forToken(baseToken).get();
|
||||
EndpointsForToken naturalViewReplicas = placement.reads.forToken(viewToken).get();
|
||||
|
||||
Optional<Replica> localReplica = Iterables.tryFind(naturalViewReplicas, Replica::isSelf).toJavaUtil();
|
||||
if (localReplica.isPresent())
|
||||
|
|
|
|||
|
|
@ -25,9 +25,7 @@ import com.google.common.base.Preconditions;
|
|||
|
||||
import org.apache.cassandra.db.Keyspace;
|
||||
import org.apache.cassandra.dht.Token;
|
||||
import org.apache.cassandra.schema.ReplicationParams;
|
||||
import org.apache.cassandra.tcm.ClusterMetadata;
|
||||
import org.apache.cassandra.tcm.ownership.DataPlacement;
|
||||
import org.apache.cassandra.tcm.ownership.VersionedEndpoints;
|
||||
|
||||
/**
|
||||
|
|
@ -166,11 +164,7 @@ public class EndpointsForToken extends Endpoints<EndpointsForToken>
|
|||
|
||||
public static VersionedEndpoints.ForToken natural(Keyspace keyspace, Token token)
|
||||
{
|
||||
ReplicationParams replication = keyspace.getMetadata().params.replication;
|
||||
DataPlacement placement = replication.isMeta()
|
||||
? ClusterMetadata.current().getCMSPlacement()
|
||||
: ClusterMetadata.current().placements.get(replication);
|
||||
return placement.reads.forToken(token);
|
||||
return ClusterMetadata.current().placement(keyspace.getMetadata().params.replication).reads.forToken(token);
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -82,13 +82,13 @@ public class MetaStrategy extends SystemStrategy
|
|||
@Override
|
||||
public EndpointsForRange calculateNaturalReplicas(Token token, ClusterMetadata metadata)
|
||||
{
|
||||
return metadata.placements.get(ReplicationParams.meta(metadata)).reads.forRange(entireRange).get();
|
||||
return metadata.placement(ReplicationParams.meta(metadata)).reads.forRange(entireRange).get();
|
||||
}
|
||||
|
||||
@Override
|
||||
public DataPlacement calculateDataPlacement(Epoch epoch, List<Range<Token>> ranges, ClusterMetadata metadata)
|
||||
{
|
||||
return metadata.placements.get(ReplicationParams.meta(metadata));
|
||||
return metadata.placement(ReplicationParams.meta(metadata));
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -239,9 +239,7 @@ public abstract class ReplicaLayout<E extends Endpoints<E>>
|
|||
{
|
||||
// todo deduplicate so that "pending" contains "read - write",
|
||||
// which is a hack until we revisit how consistency level handles pending
|
||||
DataPlacement dataPlacement = ks.params.replication.isMeta()
|
||||
? metadata.getCMSPlacement()
|
||||
: metadata.placements.get(ks.params.replication);
|
||||
DataPlacement dataPlacement = metadata.placement(ks.params.replication);
|
||||
natural = forNonLocalStrategyTokenRead(dataPlacement, token);
|
||||
// perf optimization to avoid double endpoints search and filtering for a typical case
|
||||
// DataPlacement constructor does a deduplication of reads/writes, so we can use cheap == comparision here
|
||||
|
|
@ -396,15 +394,12 @@ public abstract class ReplicaLayout<E extends Endpoints<E>>
|
|||
|
||||
static EndpointsForRange forNonLocalStategyRangeRead(ClusterMetadata metadata, KeyspaceMetadata keyspace, AbstractBounds<PartitionPosition> range)
|
||||
{
|
||||
DataPlacement placement = keyspace.params.replication.isMeta()
|
||||
? metadata.getCMSPlacement()
|
||||
: metadata.placements.get(keyspace.params.replication);
|
||||
return placement.reads.forRange(range.right.getToken()).get();
|
||||
return metadata.placement(keyspace.params.replication).reads.forRange(range.right.getToken()).get();
|
||||
}
|
||||
|
||||
public static EndpointsForToken forNonLocalStrategyTokenRead(ClusterMetadata metadata, KeyspaceMetadata keyspace, Token token)
|
||||
{
|
||||
return forNonLocalStrategyTokenRead(metadata.placements.get(keyspace.params.replication), token);
|
||||
return forNonLocalStrategyTokenRead(metadata.placement(keyspace.params.replication), token);
|
||||
}
|
||||
|
||||
public static EndpointsForToken forNonLocalStrategyTokenRead(DataPlacement dataPlacement, Token token)
|
||||
|
|
|
|||
|
|
@ -234,7 +234,7 @@ public class ReplicaPlans
|
|||
NodeProximity proximity = DatabaseDescriptor.getNodeProximity();
|
||||
AbstractReplicationStrategy replicationStrategy = keyspace.getReplicationStrategy();
|
||||
|
||||
EndpointsForToken replicas = metadata.placements.get(keyspace.getMetadata().params.replication).reads.forToken(key.getToken()).get();
|
||||
EndpointsForToken replicas = metadata.placement(keyspace.getMetadata().params.replication).reads.forToken(key.getToken()).get();
|
||||
|
||||
// CASSANDRA-13043: filter out those endpoints not accepting clients yet, maybe because still bootstrapping
|
||||
replicas = replicas.filter(replica -> StorageService.instance.isRpcReady(replica.endpoint()));
|
||||
|
|
|
|||
|
|
@ -1213,7 +1213,7 @@ public class ActiveRepairService implements IEndpointStateChangeSubscriber, IFai
|
|||
// are based on the system partitioner
|
||||
EndpointsForRange endpoints = replication.isMeta()
|
||||
? ClusterMetadata.current().fullCMSMembersAsReplicas()
|
||||
: ClusterMetadata.current().placements.get(replication).reads.forRange(range).get();
|
||||
: ClusterMetadata.current().placement(replication).reads.forRange(range).get();
|
||||
|
||||
Set<InetAddressAndPort> liveEndpoints = endpoints.filter(FailureDetector.isReplicaAlive).endpoints();
|
||||
if (!PaxosRepair.hasSufficientLiveNodesForTopologyChange(keyspace, range, liveEndpoints))
|
||||
|
|
|
|||
|
|
@ -232,7 +232,7 @@ public class Rebuild
|
|||
private static MovementMap movementMap(ClusterMetadata metadata, String keyspace, String tokens)
|
||||
{
|
||||
MovementMap.Builder movementMapBuilder = MovementMap.builder();
|
||||
DataPlacements placements = metadata.placements;
|
||||
DataPlacements placements = metadata.placements();
|
||||
if (keyspace == null)
|
||||
{
|
||||
placements.forEach((params, placement) -> movementMapBuilder.put(params, addMovementsForParams(placement, null)));
|
||||
|
|
|
|||
|
|
@ -2119,7 +2119,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
|
|||
{
|
||||
if (keyspaceMetadata.params.replication.isMeta())
|
||||
{
|
||||
DataPlacement placement = metadata.placements.get(keyspaceMetadata.params.replication);
|
||||
DataPlacement placement = metadata.placement(keyspaceMetadata.params.replication);
|
||||
// May be empty if mid-upgrade and CMS is not yet initialized
|
||||
if (!placement.reads.isEmpty())
|
||||
rangeToEndpointMap.put(MetaStrategy.entireRange, placement.reads.forRange(MetaStrategy.entireRange).get());
|
||||
|
|
@ -2129,8 +2129,8 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
|
|||
TokenMap tokenMap = metadata.tokenMap;
|
||||
for (Range<Token> range : ranges)
|
||||
{
|
||||
Token token = tokenMap.nextToken(tokenMap.tokens(), range.right.getToken());
|
||||
rangeToEndpointMap.put(range, metadata.placements.get(keyspaceMetadata.params.replication)
|
||||
Token token = TokenMap.nextToken(tokenMap.tokens(), range.right.getToken());
|
||||
rangeToEndpointMap.put(range, metadata.placement(keyspaceMetadata.params.replication)
|
||||
.reads.forRange(token).get());
|
||||
}
|
||||
}
|
||||
|
|
@ -3447,14 +3447,14 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
|
|||
token = MetaStrategy.partitioner.getToken(key);
|
||||
else
|
||||
token = metadata.partitioner.getToken(key);
|
||||
return metadata.placements.get(keyspaceMetadata.params.replication).reads.forToken(token).get();
|
||||
return metadata.placement(keyspaceMetadata.params.replication).reads.forToken(token).get();
|
||||
}
|
||||
|
||||
public boolean isEndpointValidForWrite(String keyspace, Token token)
|
||||
{
|
||||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
KeyspaceMetadata keyspaceMetadata = metadata.schema.getKeyspaces().getNullable(keyspace);
|
||||
return keyspaceMetadata != null && metadata.placements.get(keyspaceMetadata.params.replication).writes.forToken(token).get().containsSelf();
|
||||
return keyspaceMetadata != null && metadata.placement(keyspaceMetadata.params.replication).writes.forToken(token).get().containsSelf();
|
||||
}
|
||||
|
||||
public void setLoggingLevel(String classQualifier, String rawLevel) throws Exception
|
||||
|
|
@ -4262,7 +4262,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
|
|||
if (replicationParams.isMeta())
|
||||
{
|
||||
LinkedHashMap<InetAddressAndPort, Float> ownership = Maps.newLinkedHashMap();
|
||||
metadata.placements.get(replicationParams).writes.byEndpoint().flattenValues().forEach((r) -> {
|
||||
metadata.placement(replicationParams).writes.byEndpoint().flattenValues().forEach((r) -> {
|
||||
ownership.put(r.endpoint(), 1.0f);
|
||||
});
|
||||
return ownership;
|
||||
|
|
|
|||
|
|
@ -277,7 +277,7 @@ public class Paxos
|
|||
final Token token = table.partitioner == MetaStrategy.partitioner ? MetaStrategy.entireRange.right : key.getToken();
|
||||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
Keyspace keyspace = Keyspace.open(table.keyspace);
|
||||
DataPlacement placement = metadata.placements.get(keyspace.getMetadata().params.replication);
|
||||
DataPlacement placement = metadata.placement(keyspace.getMetadata().params.replication);
|
||||
Epoch epoch = placement.writes.forToken(token).lastModified();
|
||||
ForTokenWrite electorate = forTokenWriteLiveAndDown(metadata, keyspace, token);
|
||||
if (consistency == LOCAL_SERIAL)
|
||||
|
|
|
|||
|
|
@ -586,7 +586,7 @@ public class PaxosRepair extends AbstractPaxosRepair
|
|||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
Collection<InetAddressAndPort> allEndpoints = replication.isMeta()
|
||||
? metadata.fullCMSMembers()
|
||||
: metadata.placements.get(replication).reads.forRange(range).endpoints();
|
||||
: metadata.placement(replication).reads.forRange(range).endpoints();
|
||||
return hasSufficientLiveNodesForTopologyChange(allEndpoints,
|
||||
liveEndpoints,
|
||||
ep -> metadata.locator.location(ep).datacenter,
|
||||
|
|
|
|||
|
|
@ -50,7 +50,7 @@ public class DataMovementVerbHandler implements IVerbHandler<DataMovement>
|
|||
StreamPlan streamPlan = new StreamPlan(StreamOperation.fromString(message.payload.streamOperation));
|
||||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
Schema.instance.getNonLocalStrategyKeyspaces().stream().forEach((ksm) -> {
|
||||
if (metadata.placements.get(ksm.params.replication).writes.byEndpoint().keySet().size() <= 1)
|
||||
if (metadata.placement(ksm.params.replication).writes.byEndpoint().keySet().size() <= 1)
|
||||
return;
|
||||
|
||||
message.payload.movements.get(ksm.params.replication).asMap().forEach((local, endpoints) -> {
|
||||
|
|
|
|||
|
|
@ -319,6 +319,16 @@ public class ClusterMetadata
|
|||
}
|
||||
}
|
||||
|
||||
public DataPlacement placement(ReplicationParams params)
|
||||
{
|
||||
return params.isMeta() ? cmsDataPlacement : placements.get(params);
|
||||
}
|
||||
|
||||
public DataPlacements placements()
|
||||
{
|
||||
return placements;
|
||||
}
|
||||
|
||||
public Transformer transformer()
|
||||
{
|
||||
return new Transformer(this, this.nextEpoch());
|
||||
|
|
@ -484,8 +494,8 @@ public class ClusterMetadata
|
|||
// TODO Remove this as it isn't really an equivalent to the previous concept of pending ranges
|
||||
public boolean hasPendingRangesFor(KeyspaceMetadata ksm, Token token)
|
||||
{
|
||||
ReplicaGroups writes = placements.get(ksm.params.replication).writes;
|
||||
ReplicaGroups reads = placements.get(ksm.params.replication).reads;
|
||||
ReplicaGroups writes = placement(ksm.params.replication).writes;
|
||||
ReplicaGroups reads = placement(ksm.params.replication).reads;
|
||||
if (ksm.params.replication.isMeta())
|
||||
return !reads.equals(writes);
|
||||
return !reads.forToken(token).equals(writes.forToken(token));
|
||||
|
|
@ -494,8 +504,8 @@ public class ClusterMetadata
|
|||
// TODO Remove this as it isn't really an equivalent to the previous concept of pending ranges
|
||||
public boolean hasPendingRangesFor(KeyspaceMetadata ksm, InetAddressAndPort endpoint)
|
||||
{
|
||||
ReplicaGroups writes = placements.get(ksm.params.replication).writes;
|
||||
ReplicaGroups reads = placements.get(ksm.params.replication).reads;
|
||||
ReplicaGroups writes = placement(ksm.params.replication).writes;
|
||||
ReplicaGroups reads = placement(ksm.params.replication).reads;
|
||||
return !writes.byEndpoint().get(endpoint).equals(reads.byEndpoint().get(endpoint));
|
||||
}
|
||||
|
||||
|
|
@ -506,22 +516,22 @@ public class ClusterMetadata
|
|||
|
||||
public RangesAtEndpoint writeRanges(KeyspaceMetadata metadata, InetAddressAndPort peer)
|
||||
{
|
||||
return placements.get(metadata.params.replication).writes.byEndpoint().get(peer);
|
||||
return placement(metadata.params.replication).writes.byEndpoint().get(peer);
|
||||
}
|
||||
|
||||
// TODO Remove this as it isn't really an equivalent to the previous concept of pending ranges
|
||||
public Map<Range<Token>, VersionedEndpoints.ForRange> pendingRanges(KeyspaceMetadata metadata)
|
||||
{
|
||||
Map<Range<Token>, VersionedEndpoints.ForRange> map = new HashMap<>();
|
||||
ReplicaGroups writes = placements.get(metadata.params.replication).writes;
|
||||
ReplicaGroups reads = placements.get(metadata.params.replication).reads;
|
||||
ReplicaGroups writes = placement(metadata.params.replication).writes;
|
||||
ReplicaGroups reads = placement(metadata.params.replication).reads;
|
||||
|
||||
// first, pending ranges as the result of range splitting or merging
|
||||
// i.e. new ranges being created through join/leave
|
||||
List<Range<Token>> pending = new ArrayList<>(writes.ranges());
|
||||
pending.removeAll(reads.ranges());
|
||||
for (Range<Token> p : pending)
|
||||
map.put(p, placements.get(metadata.params.replication).writes.forRange(p));
|
||||
map.put(p, placement(metadata.params.replication).writes.forRange(p));
|
||||
|
||||
// next, ranges where the ranges themselves are not changing, but the replicas are
|
||||
// i.e. replacement or RF increase
|
||||
|
|
@ -538,8 +548,8 @@ public class ClusterMetadata
|
|||
// TODO Remove this as it isn't really an equivalent to the previous concept of pending endpoints
|
||||
public VersionedEndpoints.ForToken pendingEndpointsFor(KeyspaceMetadata metadata, Token t)
|
||||
{
|
||||
VersionedEndpoints.ForToken writeEndpoints = placements.get(metadata.params.replication).writes.forToken(t);
|
||||
VersionedEndpoints.ForToken readEndpoints = placements.get(metadata.params.replication).reads.forToken(t);
|
||||
VersionedEndpoints.ForToken writeEndpoints = placement(metadata.params.replication).writes.forToken(t);
|
||||
VersionedEndpoints.ForToken readEndpoints = placement(metadata.params.replication).reads.forToken(t);
|
||||
EndpointsForToken.Builder endpointsForToken = writeEndpoints.get().newBuilder(writeEndpoints.size() - readEndpoints.size());
|
||||
|
||||
for (Replica writeReplica : writeEndpoints.get())
|
||||
|
|
|
|||
|
|
@ -43,7 +43,6 @@ import org.apache.cassandra.net.MessagingService;
|
|||
import org.apache.cassandra.net.RequestCallbackWithFailure;
|
||||
import org.apache.cassandra.net.Verb;
|
||||
import org.apache.cassandra.schema.DistributedMetadataLogKeyspace;
|
||||
import org.apache.cassandra.schema.ReplicationParams;
|
||||
import org.apache.cassandra.tcm.log.Entry;
|
||||
import org.apache.cassandra.tcm.log.LocalLog;
|
||||
import org.apache.cassandra.tcm.log.LogState;
|
||||
|
|
|
|||
|
|
@ -33,8 +33,8 @@ public class PlacementsChangeListener implements ChangeListener
|
|||
|
||||
private boolean shouldInvalidate(ClusterMetadata prev, ClusterMetadata next)
|
||||
{
|
||||
if (!prev.placements.lastModified().equals(next.placements.lastModified()) &&
|
||||
!prev.placements.equivalentTo(next.placements)) // <- todo should we update lastModified if the result is the same?
|
||||
if (!prev.placements().lastModified().equals(next.placements().lastModified()) &&
|
||||
!prev.placements().equivalentTo(next.placements())) // <- todo should we update lastModified if the result is the same?
|
||||
return true;
|
||||
|
||||
if (prev.schema.getKeyspaces().size() != next.schema.getKeyspaces().size())
|
||||
|
|
|
|||
|
|
@ -335,7 +335,7 @@ public class UniformRangePlacement implements PlacementProvider
|
|||
}
|
||||
else if (params.isLocal())
|
||||
{
|
||||
placements.put(params, metadata.placements.get(params));
|
||||
placements.put(params, metadata.placement(params));
|
||||
}
|
||||
else
|
||||
{
|
||||
|
|
|
|||
|
|
@ -326,7 +326,7 @@ public class BootstrapAndJoin extends MultiStepOperation<Epoch>
|
|||
@Override
|
||||
public ClusterMetadata.Transformer cancel(ClusterMetadata metadata)
|
||||
{
|
||||
DataPlacements placements = metadata.placements;
|
||||
DataPlacements placements = metadata.placements();
|
||||
switch (next)
|
||||
{
|
||||
// need to undo MID_JOIN and START_JOIN, then merge the ranges split by PrepareJoin
|
||||
|
|
@ -357,7 +357,7 @@ public class BootstrapAndJoin extends MultiStepOperation<Epoch>
|
|||
@VisibleForTesting
|
||||
public Pair<MovementMap, MovementMap> getMovementMaps(ClusterMetadata metadata)
|
||||
{
|
||||
MovementMap movementMap = movementMap(metadata.directory.endpoint(startJoin.nodeId()), metadata.placements, startJoin.delta());
|
||||
MovementMap movementMap = movementMap(metadata.directory.endpoint(startJoin.nodeId()), metadata.placements(), startJoin.delta());
|
||||
MovementMap strictMovementMap = toStrict(movementMap, finishJoin.delta());
|
||||
return Pair.create(movementMap, strictMovementMap);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -318,7 +318,7 @@ public class BootstrapAndReplace extends MultiStepOperation<Epoch>
|
|||
@Override
|
||||
public ClusterMetadata.Transformer cancel(ClusterMetadata metadata)
|
||||
{
|
||||
DataPlacements placements = metadata.placements;
|
||||
DataPlacements placements = metadata.placements();
|
||||
switch (next)
|
||||
{
|
||||
// need to undo MID_REPLACE and START_REPLACE, but PREPARE_REPLACE doesn't affect placements
|
||||
|
|
@ -355,7 +355,7 @@ public class BootstrapAndReplace extends MultiStepOperation<Epoch>
|
|||
private static MovementMap movementMap(InetAddressAndPort beingReplaced, PlacementDeltas startDelta)
|
||||
{
|
||||
MovementMap.Builder movementMapBuilder = MovementMap.builder();
|
||||
DataPlacements placements = ClusterMetadata.current().placements;
|
||||
DataPlacements placements = ClusterMetadata.current().placements();
|
||||
startDelta.forEach((params, delta) -> {
|
||||
EndpointsByReplica.Builder movements = new EndpointsByReplica.Builder();
|
||||
DataPlacement originalPlacements = placements.get(params);
|
||||
|
|
|
|||
|
|
@ -321,7 +321,7 @@ public class Move extends MultiStepOperation<Epoch>
|
|||
StreamPlan streamPlan = new StreamPlan(StreamOperation.RELOCATION);
|
||||
Keyspaces keyspaces = Schema.instance.getNonLocalStrategyKeyspaces();
|
||||
Map<ReplicationParams, EndpointsByReplica> movementMap = movementMap(FailureDetector.instance,
|
||||
metadata.placements,
|
||||
metadata.placements(),
|
||||
toSplitRanges,
|
||||
startMove.delta(),
|
||||
midMove.delta(),
|
||||
|
|
@ -430,7 +430,7 @@ public class Move extends MultiStepOperation<Epoch>
|
|||
@Override
|
||||
public ClusterMetadata.Transformer cancel(ClusterMetadata metadata)
|
||||
{
|
||||
DataPlacements placements = metadata.placements;
|
||||
DataPlacements placements = metadata.placements();
|
||||
|
||||
switch (next)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -58,6 +58,7 @@ import org.apache.cassandra.tcm.Epoch;
|
|||
import org.apache.cassandra.tcm.Retry;
|
||||
import org.apache.cassandra.tcm.membership.Directory;
|
||||
import org.apache.cassandra.tcm.membership.Location;
|
||||
import org.apache.cassandra.tcm.ownership.DataPlacement;
|
||||
import org.apache.cassandra.utils.Clock;
|
||||
import org.apache.cassandra.utils.concurrent.AsyncPromise;
|
||||
|
||||
|
|
@ -170,8 +171,9 @@ public class ProgressBarrier
|
|||
Set<Range<Token>> ranges = e.getValue();
|
||||
for (Range<Token> range : ranges)
|
||||
{
|
||||
EndpointsForRange writes = metadata.placements.get(params).writes.matchRange(range).get().filter(r -> filter.test(r.endpoint()));
|
||||
EndpointsForRange reads = metadata.placements.get(params).reads.matchRange(range).get().filter(r -> filter.test(r.endpoint()));
|
||||
DataPlacement placement = metadata.placement(params);
|
||||
EndpointsForRange writes = placement.writes.matchRange(range).get().filter(r -> filter.test(r.endpoint()));
|
||||
EndpointsForRange reads = placement.reads.matchRange(range).get().filter(r -> filter.test(r.endpoint()));
|
||||
// Affected ranges can contain ranges which are the results of merging or splitting and may not exist
|
||||
// as keys in the existing ReplicaGroups. As such, no replicas will be found for these ranges and so no
|
||||
// WaitFor is necessary.
|
||||
|
|
|
|||
|
|
@ -121,7 +121,7 @@ public class RemoveNodeStreams implements LeaveStreams
|
|||
RangesByEndpoint startWriteAdditions = startDelta.get(params).writes.additions;
|
||||
RangesByEndpoint startWriteRemovals = startDelta.get(params).writes.removals;
|
||||
// find current placements from the metadata, we need to stream from replicas that are not changed and are therefore not in the deltas
|
||||
ReplicaGroups currentPlacements = metadata.placements.get(params).reads;
|
||||
ReplicaGroups currentPlacements = metadata.placement(params).reads;
|
||||
startWriteAdditions.flattenValues()
|
||||
.forEach(newReplica -> {
|
||||
EndpointsForRange.Builder candidateBuilder = new EndpointsForRange.Builder(newReplica.range());
|
||||
|
|
|
|||
|
|
@ -44,7 +44,7 @@ public class ReplaceSameAddress
|
|||
{
|
||||
MovementMap.Builder builder = MovementMap.builder();
|
||||
InetAddressAndPort addr = metadata.directory.endpoint(nodeId);
|
||||
metadata.placements.forEach((params, placement) -> {
|
||||
metadata.placements().forEach((params, placement) -> {
|
||||
EndpointsByReplica.Builder sources = new EndpointsByReplica.Builder();
|
||||
placement.reads.byEndpoint().get(addr).forEach(destination -> {
|
||||
placement.reads.forRange(destination.range()).forEach(potentialSource -> {
|
||||
|
|
|
|||
|
|
@ -253,7 +253,7 @@ public class UnbootstrapAndLeave extends MultiStepOperation<Epoch>
|
|||
@Override
|
||||
public ClusterMetadata.Transformer cancel(ClusterMetadata metadata)
|
||||
{
|
||||
DataPlacements placements = metadata.placements;
|
||||
DataPlacements placements = metadata.placements();
|
||||
switch (next)
|
||||
{
|
||||
// need to undo MID_LEAVE and START_LEAVE, but PrepareLeave doesn't affect placement
|
||||
|
|
|
|||
|
|
@ -244,7 +244,7 @@ public class AlterSchema implements Transformation
|
|||
|
||||
DataPlacements.Builder newPlacementsBuilder = DataPlacements.builder(calculatedPlacements.size());
|
||||
calculatedPlacements.forEach((params, newPlacement) -> {
|
||||
DataPlacement previousPlacement = prev.placements.get(params);
|
||||
DataPlacement previousPlacement = prev.placement(params);
|
||||
// Preserve placement versioning that has resulted from natural application where possible
|
||||
if (previousPlacement.equivalentTo(newPlacement))
|
||||
newPlacementsBuilder.with(params, previousPlacement);
|
||||
|
|
|
|||
|
|
@ -130,11 +130,11 @@ public class AlterTopology implements Transformation
|
|||
for (Map.Entry<NodeId, Location> update : updates.entrySet())
|
||||
updated = updated.withUpdatedRackAndDc(update.getKey(), update.getValue());
|
||||
ClusterMetadata proposed = prev.transformer().with(updated).build().metadata;
|
||||
DataPlacements proposedPlacements = placementProvider.calculatePlacements(prev.placements.lastModified(),
|
||||
DataPlacements proposedPlacements = placementProvider.calculatePlacements(prev.placements().lastModified(),
|
||||
proposed.tokenMap.toRanges(),
|
||||
proposed,
|
||||
proposed.schema.getKeyspaces());
|
||||
if (!proposedPlacements.equivalentTo(prev.placements))
|
||||
if (!proposedPlacements.equivalentTo(prev.placements()))
|
||||
{
|
||||
logger.info("Rejecting topology modifications which would materially change data placements: {}", updates);
|
||||
return new Rejected(INVALID, "Proposed updates modify data placements, violating consistency guarantees");
|
||||
|
|
|
|||
|
|
@ -82,7 +82,7 @@ public abstract class ApplyPlacementDeltas implements Transformation
|
|||
ClusterMetadata.Transformer next = prev.transformer();
|
||||
|
||||
if (!delta.isEmpty())
|
||||
next = next.with(delta.apply(prev.nextEpoch(), prev.placements));
|
||||
next = next.with(delta.apply(prev.nextEpoch(), prev.placements()));
|
||||
|
||||
next = transform(prev, next);
|
||||
|
||||
|
|
|
|||
|
|
@ -89,7 +89,7 @@ import static org.apache.cassandra.exceptions.ExceptionCode.INVALID;
|
|||
*/
|
||||
public class PrepareJoin implements Transformation
|
||||
{
|
||||
public static final Serializer<PrepareJoin> serializer = new Serializer<PrepareJoin>()
|
||||
public static final Serializer<PrepareJoin> serializer = new Serializer<>()
|
||||
{
|
||||
public PrepareJoin construct(NodeId nodeId, Set<Token> tokens, PlacementProvider placementProvider, boolean joinTokenRing, boolean streamData)
|
||||
{
|
||||
|
|
@ -168,10 +168,10 @@ public class PrepareJoin implements Transformation
|
|||
startJoin, midJoin, finishJoin,
|
||||
joinTokenRing, streamData);
|
||||
if (!prev.tokenMap.isEmpty())
|
||||
assertPreExistingWriteReplica(prev.placements, transitionPlan);
|
||||
assertPreExistingWriteReplica(prev.placements(), transitionPlan);
|
||||
|
||||
LockedRanges newLockedRanges = prev.lockedRanges.lock(lockKey, rangesToLock);
|
||||
DataPlacements startingPlacements = transitionPlan.toSplit.apply(prev.nextEpoch(), prev.placements);
|
||||
DataPlacements startingPlacements = transitionPlan.toSplit.apply(prev.nextEpoch(), prev.placements());
|
||||
ClusterMetadata.Transformer proposed = prev.transformer()
|
||||
.with(newLockedRanges)
|
||||
.with(startingPlacements)
|
||||
|
|
|
|||
|
|
@ -115,7 +115,7 @@ public class PrepareLeave implements Transformation
|
|||
PlacementDeltas startDelta = transitionPlan.addToWrites();
|
||||
PlacementDeltas midDelta = transitionPlan.moveReads();
|
||||
PlacementDeltas finishDelta = transitionPlan.removeFromWrites();
|
||||
transitionPlan.assertPreExistingWriteReplica(prev.placements);
|
||||
transitionPlan.assertPreExistingWriteReplica(prev.placements());
|
||||
|
||||
LockedRanges.Key unlockKey = LockedRanges.keyFor(proposed.epoch);
|
||||
|
||||
|
|
|
|||
|
|
@ -109,7 +109,7 @@ public class PrepareMove implements Transformation
|
|||
StartMove startMove = new StartMove(nodeId, transitionPlan.addToWrites(), lockKey);
|
||||
MidMove midMove = new MidMove(nodeId, transitionPlan.moveReads(), lockKey);
|
||||
FinishMove finishMove = new FinishMove(nodeId, tokens, transitionPlan.removeFromWrites(), lockKey);
|
||||
transitionPlan.assertPreExistingWriteReplica(prev.placements);
|
||||
transitionPlan.assertPreExistingWriteReplica(prev.placements());
|
||||
|
||||
Move sequence = Move.newSequence(prev.nextEpoch(),
|
||||
lockKey,
|
||||
|
|
@ -123,7 +123,7 @@ public class PrepareMove implements Transformation
|
|||
return Transformation.success(prev.transformer()
|
||||
.withNodeState(nodeId, NodeState.MOVING)
|
||||
.with(prev.lockedRanges.lock(lockKey, rangesToLock))
|
||||
.with(transitionPlan.toSplit.apply(prev.nextEpoch(), prev.placements))
|
||||
.with(transitionPlan.toSplit.apply(prev.nextEpoch(), prev.placements()))
|
||||
.with(prev.inProgressSequences.with(nodeId, sequence)),
|
||||
rangesToLock);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -109,7 +109,7 @@ public class PrepareReplace implements Transformation
|
|||
StartReplace start = new StartReplace(replaced, replacement, transitionPlan.addToWrites(), unlockKey);
|
||||
MidReplace mid = new MidReplace(replaced, replacement, transitionPlan.moveReads(), unlockKey);
|
||||
FinishReplace finish = new FinishReplace(replaced, replacement, transitionPlan.removeFromWrites(), unlockKey);
|
||||
transitionPlan.assertPreExistingWriteReplica(prev.placements);
|
||||
transitionPlan.assertPreExistingWriteReplica(prev.placements());
|
||||
|
||||
Set<Token> tokens = new HashSet<>(prev.tokenMap.tokens(replaced));
|
||||
BootstrapAndReplace plan = BootstrapAndReplace.newSequence(prev.nextEpoch(),
|
||||
|
|
|
|||
|
|
@ -266,9 +266,7 @@ public abstract class PrepareCMSReconfiguration implements Transformation
|
|||
KeyspaceMetadata keyspace = prev.schema.getKeyspaceMetadata(SchemaConstants.METADATA_KEYSPACE_NAME);
|
||||
KeyspaceMetadata newKeyspace = keyspace.withSwapped(new KeyspaceParams(keyspace.params.durableWrites, replicationParams, FastPathStrategy.simple()));
|
||||
|
||||
return executeInternal(prev,
|
||||
transformer -> transformer.with(prev.placements.replaceParams(prev.nextEpoch(), ReplicationParams.meta(prev), replicationParams))
|
||||
.with(new DistributedSchema(prev.schema.getKeyspaces().withAddedOrUpdated(newKeyspace))));
|
||||
return executeInternal(prev, transformer -> transformer.with(new DistributedSchema(prev.schema.getKeyspaces().withAddedOrUpdated(newKeyspace))));
|
||||
}
|
||||
|
||||
public String toString()
|
||||
|
|
|
|||
|
|
@ -27,13 +27,11 @@ import org.apache.cassandra.dht.IPartitioner;
|
|||
import org.apache.cassandra.io.util.FileInputStreamPlus;
|
||||
import org.apache.cassandra.io.util.FileOutputStreamPlus;
|
||||
import org.apache.cassandra.locator.InetAddressAndPort;
|
||||
import org.apache.cassandra.locator.MetaStrategy;
|
||||
import org.apache.cassandra.locator.Replica;
|
||||
import org.apache.cassandra.schema.ReplicationParams;
|
||||
import org.apache.cassandra.tcm.ClusterMetadata;
|
||||
import org.apache.cassandra.tcm.ClusterMetadataService;
|
||||
import org.apache.cassandra.tcm.membership.NodeId;
|
||||
import org.apache.cassandra.tcm.membership.NodeVersion;
|
||||
import org.apache.cassandra.tcm.ownership.DataPlacement;
|
||||
import org.apache.cassandra.tcm.serialization.VerboseMetadataSerializer;
|
||||
import org.apache.cassandra.tcm.serialization.Version;
|
||||
|
||||
|
|
@ -60,9 +58,9 @@ public class TransformClusterMetadataHelper
|
|||
DatabaseDescriptor.setPartitionerUnsafe(partitioner);
|
||||
ClusterMetadataService.initializeForTools(false);
|
||||
ClusterMetadata metadata = ClusterMetadataService.deserializeClusterMetadata(sourceFile);
|
||||
System.out.println("Old CMS: " + metadata.placements.get(ReplicationParams.meta(metadata)));
|
||||
System.out.println("Old CMS: " + metadata.placement(ReplicationParams.meta(metadata)));
|
||||
metadata = makeCMS(metadata, InetAddressAndPort.getByNameUnchecked(args[1]));
|
||||
System.out.println("New CMS: " + metadata.placements.get(ReplicationParams.meta(metadata)));
|
||||
System.out.println("New CMS: " + metadata.placement(ReplicationParams.meta(metadata)));
|
||||
Path p = Files.createTempFile("clustermetadata", "dump");
|
||||
try (FileOutputStreamPlus out = new FileOutputStreamPlus(p))
|
||||
{
|
||||
|
|
@ -73,20 +71,9 @@ public class TransformClusterMetadataHelper
|
|||
|
||||
public static ClusterMetadata makeCMS(ClusterMetadata metadata, InetAddressAndPort endpoint)
|
||||
{
|
||||
ReplicationParams metaParams = ReplicationParams.meta(metadata);
|
||||
Iterable<Replica> currentReplicas = metadata.placements.get(metaParams).writes.byEndpoint().flattenValues();
|
||||
DataPlacement.Builder builder = metadata.placements.get(metaParams).unbuild();
|
||||
for (Replica replica : currentReplicas)
|
||||
{
|
||||
builder.withoutReadReplica(metadata.epoch, replica)
|
||||
.withoutWriteReplica(metadata.epoch, replica);
|
||||
}
|
||||
Replica newCMS = MetaStrategy.replica(endpoint);
|
||||
builder.withReadReplica(metadata.epoch, newCMS)
|
||||
.withWriteReplica(metadata.epoch, newCMS);
|
||||
return metadata.transformer().with(metadata.placements.unbuild().with(metaParams,
|
||||
builder.build())
|
||||
.build())
|
||||
.build().metadata;
|
||||
NodeId id = metadata.directory.peerId(endpoint);
|
||||
if (id == null)
|
||||
throw new IllegalStateException("No node id found for endpoint: " + endpoint);
|
||||
return metadata.transformer().startJoiningCMS(id).finishJoiningCMS(id).build().metadata;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -444,8 +444,8 @@ public class ClusterUtils
|
|||
for (KeyspaceMetadata keyspace : metadata.schema.getKeyspaces())
|
||||
{
|
||||
List[] placements = new List[2];
|
||||
placements[0] = metadata.placements.get(keyspace.params.replication).reads.toReplicaStringList();
|
||||
placements[1] = metadata.placements.get(keyspace.params.replication).writes.toReplicaStringList();
|
||||
placements[0] = metadata.placement(keyspace.params.replication).reads.toReplicaStringList();
|
||||
placements[1] = metadata.placement(keyspace.params.replication).writes.toReplicaStringList();
|
||||
byKeyspace.put(keyspace.name, placements);
|
||||
}
|
||||
return byKeyspace;
|
||||
|
|
@ -465,10 +465,10 @@ public class ClusterUtils
|
|||
StringBuilder builder = new StringBuilder();
|
||||
builder.append("'keyspace' { 'name':").append(keyspace.name).append("', ");
|
||||
builder.append("'reads':['");
|
||||
ReplicaGroups placement = metadata.placements.get(keyspace.params.replication).reads;
|
||||
ReplicaGroups placement = metadata.placement(keyspace.params.replication).reads;
|
||||
builder.append(byEndpoint ? placement.toStringByEndpoint() : placement.toString());
|
||||
builder.append("'], 'writes':['");
|
||||
placement = metadata.placements.get(keyspace.params.replication).writes;
|
||||
placement = metadata.placement(keyspace.params.replication).writes;
|
||||
builder.append(byEndpoint ? placement.toStringByEndpoint() : placement.toString());
|
||||
builder.append("']}");
|
||||
keyspaces.add(builder.toString());
|
||||
|
|
|
|||
|
|
@ -1077,7 +1077,7 @@ public class ClusterMetadataTestHelper
|
|||
public static VersionedEndpoints.ForToken getNaturalReplicasForToken(ClusterMetadata metadata, String keyspace, Token searchPosition)
|
||||
{
|
||||
KeyspaceMetadata keyspaceMetadata = metadata.schema.getKeyspaces().getNullable(keyspace);
|
||||
return metadata.placements.get(keyspaceMetadata.params.replication).reads.forToken(searchPosition);
|
||||
return metadata.placement(keyspaceMetadata.params.replication).reads.forToken(searchPosition);
|
||||
}
|
||||
|
||||
public static BootstrapAndJoin getBootstrapPlan(int idx)
|
||||
|
|
|
|||
|
|
@ -508,7 +508,7 @@ public class MetadataChangeSimulationTest extends CMSTestBase
|
|||
Set<NodeId> bouncing = new HashSet<>();
|
||||
Set<NodeId> replicasFromBouncedReplicaSets = new HashSet<>();
|
||||
outer:
|
||||
for (VersionedEndpoints.ForRange placements : sut.service.metadata().placements.get(rf.asKeyspaceParams().replication).writes.endpoints)
|
||||
for (VersionedEndpoints.ForRange placements : sut.service.metadata().placement(rf.asKeyspaceParams().replication).writes.endpoints)
|
||||
{
|
||||
List<NodeId> replicas = new ArrayList<>(metadata.directory.toNodeIds(placements.get().endpoints()));
|
||||
List<NodeId> bounceCandidates = new ArrayList<>();
|
||||
|
|
@ -545,8 +545,8 @@ public class MetadataChangeSimulationTest extends CMSTestBase
|
|||
ClusterMetadata actualMetadata = sut.service.metadata();
|
||||
ReplicationParams replication = actualMetadata.schema.getKeyspaces().get("test").get().params.replication;
|
||||
Assert.assertEquals(replication, sut.rf.asKeyspaceParams().replication);
|
||||
match(actualMetadata.placements.get(replication).reads, sut.rf.replicate(modelState.simulatedPlacements.nodes).asMap());
|
||||
match(actualMetadata.placements.get(replication).writes, sut.rf.replicate(modelState.simulatedPlacements.nodes).asMap());
|
||||
match(actualMetadata.placement(replication).reads, sut.rf.replicate(modelState.simulatedPlacements.nodes).asMap());
|
||||
match(actualMetadata.placement(replication).writes, sut.rf.replicate(modelState.simulatedPlacements.nodes).asMap());
|
||||
}
|
||||
|
||||
public static void validatePlacements(CMSTestBase.CMSSut sut, ModelState modelState) throws Throwable
|
||||
|
|
@ -559,7 +559,7 @@ public class MetadataChangeSimulationTest extends CMSTestBase
|
|||
Assert.assertEquals(modelState.simulatedPlacements.nodes.stream().map(Node::token).collect(Collectors.toSet()),
|
||||
actualMetadata.tokenMap.tokens().stream().map(t -> ((LongToken) t).getLongValue()).collect(Collectors.toSet()));
|
||||
|
||||
for (Map.Entry<ReplicationParams, DataPlacement> e : actualMetadata.placements.asMap().entrySet())
|
||||
for (Map.Entry<ReplicationParams, DataPlacement> e : actualMetadata.placements().asMap().entrySet())
|
||||
{
|
||||
if (!e.getKey().equals(replication))
|
||||
continue;
|
||||
|
|
@ -569,7 +569,7 @@ public class MetadataChangeSimulationTest extends CMSTestBase
|
|||
match(placement.reads, modelState.simulatedPlacements.readPlacements);
|
||||
}
|
||||
|
||||
validatePlacements(sut.partitioner, sut.rf, modelState, actualMetadata.placements);
|
||||
validatePlacements(sut.partitioner, sut.rf, modelState, actualMetadata.placements());
|
||||
}
|
||||
|
||||
public static ModelChecker.Pair<ModelState, Node> registerNewNode(ModelState state, CMSSut sut, int dcIdx, int rackIdx)
|
||||
|
|
@ -953,7 +953,7 @@ public class MetadataChangeSimulationTest extends CMSTestBase
|
|||
validatePlacements(sut, state);
|
||||
}
|
||||
// Finally verify that the predicted placements match the actual ones
|
||||
Assert.assertTrue(allSettled.equivalentTo(sut.service.metadata().placements.get(ksm.params.replication)));
|
||||
Assert.assertTrue(allSettled.equivalentTo(sut.service.metadata().placement(ksm.params.replication)));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -116,8 +116,7 @@ public class OperationalEquivalenceTest extends CMSTestBase
|
|||
withMove = ClusterMetadata.current();
|
||||
}
|
||||
|
||||
assertPlacements(simulateAndCompare(rf, equivalentNodes).placements,
|
||||
withMove.placements);
|
||||
assertPlacements(simulateAndCompare(rf, equivalentNodes).placements(), withMove.placements());
|
||||
}
|
||||
|
||||
private static ClusterMetadata simulateAndCompare(ReplicationFactor rf, List<Node> nodes) throws Exception
|
||||
|
|
|
|||
|
|
@ -85,7 +85,7 @@ public class ReconfigureCMSTest extends FuzzTestBase
|
|||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
assertEquals(5, metadata.fullCMSMembers().size());
|
||||
assertEquals(ReplicationParams.simpleMeta(5, metadata.directory.knownDatacenters()),
|
||||
metadata.placements.keys().stream().filter(ReplicationParams::isMeta).findFirst().get());
|
||||
metadata.placements().keys().stream().filter(ReplicationParams::isMeta).findFirst().get());
|
||||
});
|
||||
cluster.stream().forEach(i -> {
|
||||
Assert.assertTrue(i.executeInternal(String.format("SELECT * FROM %s.%s", SchemaConstants.METADATA_KEYSPACE_NAME, DistributedMetadataLogKeyspace.TABLE_NAME)).length > 0);
|
||||
|
|
@ -96,7 +96,7 @@ public class ReconfigureCMSTest extends FuzzTestBase
|
|||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
assertEquals(1, metadata.fullCMSMembers().size());
|
||||
assertEquals(ReplicationParams.simpleMeta(1, metadata.directory.knownDatacenters()),
|
||||
metadata.placements.keys().stream().filter(ReplicationParams::isMeta).findFirst().get());
|
||||
metadata.placements().keys().stream().filter(ReplicationParams::isMeta).findFirst().get());
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
@ -136,7 +136,7 @@ public class ReconfigureCMSTest extends FuzzTestBase
|
|||
Assert.assertNull(metadata.inProgressSequences.get(ReconfigureCMS.SequenceKey.instance));
|
||||
assertEquals(2, metadata.fullCMSMembers().size());
|
||||
ReplicationParams params = ReplicationParams.meta(metadata);
|
||||
DataPlacement placements = metadata.placements.get(params);
|
||||
DataPlacement placements = metadata.placements().get(params);
|
||||
assertTrue(placements.reads.equivalentTo(placements.writes));
|
||||
assertEquals(metadata.fullCMSMembers().size(), Integer.parseInt(params.asMap().get("dc0")));
|
||||
});
|
||||
|
|
@ -161,7 +161,7 @@ public class ReconfigureCMSTest extends FuzzTestBase
|
|||
Assert.assertNull(metadata.inProgressSequences.get(ReconfigureCMS.SequenceKey.instance));
|
||||
Assert.assertTrue(metadata.fullCMSMembers().contains(FBUtilities.getBroadcastAddressAndPort()));
|
||||
assertEquals(3, metadata.fullCMSMembers().size());
|
||||
DataPlacement placements = metadata.placements.get(ReplicationParams.meta(metadata));
|
||||
DataPlacement placements = metadata.placements().get(ReplicationParams.meta(metadata));
|
||||
Assert.assertTrue(placements.reads.equivalentTo(placements.writes));
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -146,12 +146,12 @@ public class ResumableStartupTest extends FuzzTestBase
|
|||
KeyspaceMetadata ksm = metadata.schema.getKeyspaceMetadata(keyspace);
|
||||
boolean isWriteReplica = false;
|
||||
boolean isReadReplica = false;
|
||||
for (InetAddressAndPort readReplica : metadata.placements.get(ksm.params.replication).reads.byEndpoint().keySet())
|
||||
for (InetAddressAndPort readReplica : metadata.placement(ksm.params.replication).reads.byEndpoint().keySet())
|
||||
{
|
||||
if (readReplica.getHostAddressAndPort().equals(newAddress))
|
||||
isReadReplica = true;
|
||||
}
|
||||
for (InetAddressAndPort writeReplica : metadata.placements.get(ksm.params.replication).writes.byEndpoint().keySet())
|
||||
for (InetAddressAndPort writeReplica : metadata.placement(ksm.params.replication).writes.byEndpoint().keySet())
|
||||
{
|
||||
if (writeReplica.getHostAddressAndPort().equals(newAddress))
|
||||
isWriteReplica = true;
|
||||
|
|
|
|||
|
|
@ -177,8 +177,8 @@ public abstract class SimulatedOperation
|
|||
sutActions.next();
|
||||
ClusterMetadata m2 = ClusterMetadata.current();
|
||||
|
||||
Map<Range<Token>, VersionedEndpoints.ForRange> after = m2.placements.get(simulatedState.rf.asKeyspaceParams().replication).reads.asMap();
|
||||
m1.placements.get(simulatedState.rf.asKeyspaceParams().replication).reads.forEach((k, beforePlacements) -> {
|
||||
Map<Range<Token>, VersionedEndpoints.ForRange> after = m2.placement(simulatedState.rf.asKeyspaceParams().replication).reads.asMap();
|
||||
m1.placement(simulatedState.rf.asKeyspaceParams().replication).reads.forEach((k, beforePlacements) -> {
|
||||
if (after.containsKey(k))
|
||||
{
|
||||
VersionedEndpoints.ForRange afterPlacements = after.get(k);
|
||||
|
|
|
|||
|
|
@ -69,7 +69,7 @@ public class SnapshotTest extends TestBaseImpl
|
|||
ClusterMetadata before = ClusterMetadata.current();
|
||||
ClusterMetadata after = ClusterMetadataService.instance().triggerSnapshot();
|
||||
ClusterMetadata serialized = ClusterMetadataService.instance().snapshotManager().getSnapshot(after.epoch);
|
||||
assertEquals(before.placements, serialized.placements);
|
||||
assertEquals(before.placements(), serialized.placements());
|
||||
assertEquals(before.tokenMap, serialized.tokenMap);
|
||||
assertEquals(before.directory, serialized.directory);
|
||||
assertEquals(before.schema, serialized.schema);
|
||||
|
|
|
|||
|
|
@ -54,7 +54,7 @@ public class RangeVersioningTest extends FuzzTestBase
|
|||
for (int i = 1; i <= 4; i++)
|
||||
{
|
||||
Epoch smallestSeen = null;
|
||||
for (VersionedEndpoints.ForRange fr : metadata.placements.get(ReplicationParams.simple(i)).writes.endpoints)
|
||||
for (VersionedEndpoints.ForRange fr : metadata.placement(ReplicationParams.simple(i)).writes.endpoints)
|
||||
{
|
||||
if (smallestSeen == null || fr.lastModified().isBefore(smallestSeen))
|
||||
smallestSeen = fr.lastModified();
|
||||
|
|
|
|||
|
|
@ -60,7 +60,7 @@ public class ClusterMetadataUpgradeAssassinateTest extends UpgradeTestBase
|
|||
((IInvokableInstance) i).runOnInstance(() -> {
|
||||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
InetAddressAndPort ep = InetAddressAndPort.getByNameUnchecked(host);
|
||||
metadata.placements.asMap().forEach((key, value) -> {
|
||||
metadata.placements().forEach((key, value) -> {
|
||||
if (key.isMeta())
|
||||
return;
|
||||
boolean existsInPlacements = Streams.concat(value.reads.endpoints.stream(),
|
||||
|
|
|
|||
|
|
@ -945,7 +945,7 @@ public abstract class TopologyMixupTestBase<S extends TopologyMixupTestBase.Sche
|
|||
{
|
||||
return inst.callOnInstance(() -> {
|
||||
ClusterMetadata current = ClusterMetadata.current();
|
||||
Set<InetAddressAndPort> members = current.placements.get(ReplicationParams.meta(current)).writes.byEndpoint().keySet();
|
||||
Set<InetAddressAndPort> members = current.placement(ReplicationParams.meta(current)).writes.byEndpoint().keySet();
|
||||
// Why not just use 'current.fullCMSMembers()'? That uses the "read" replicas, so "could" have less endpoints
|
||||
// It would be more consistent to use fullCMSMembers but thought process is knowing the full set is better
|
||||
// than the coordination set.
|
||||
|
|
|
|||
|
|
@ -78,7 +78,7 @@ class OnClusterReplace extends OnClusterChangeTopology
|
|||
List<Map.Entry<String, String>> repairRanges = actions.cluster.get(leaving).unsafeApplyOnThisThread(
|
||||
(String keyspaceName) -> {
|
||||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
return metadata.placements.get(metadata.schema.getKeyspace(keyspaceName).getMetadata().params.replication)
|
||||
return metadata.placement(metadata.schema.getKeyspace(keyspaceName).getMetadata().params.replication)
|
||||
.writes.ranges()
|
||||
.stream()
|
||||
.map(OnClusterReplace::toStringEntry)
|
||||
|
|
@ -93,7 +93,7 @@ class OnClusterReplace extends OnClusterChangeTopology
|
|||
(String keyspaceName, String tk) -> {
|
||||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
KeyspaceMetadata keyspaceMetadata = metadata.schema.getKeyspaces().getNullable(keyspaceName);
|
||||
return metadata.placements.get(keyspaceMetadata.params.replication).reads
|
||||
return metadata.placement(keyspaceMetadata.params.replication).reads
|
||||
.forToken(Utils.parseToken(tk))
|
||||
.get()
|
||||
.stream().map(Replica::endpoint)
|
||||
|
|
|
|||
|
|
@ -372,6 +372,6 @@ public class SimpleStrategyTest extends CassandraTestBase
|
|||
ReplicationParams replicationParams,
|
||||
Token token)
|
||||
{
|
||||
return metadata.placements.get(replicationParams).writes.forToken(token).get();
|
||||
return metadata.placement(replicationParams).writes.forToken(token).get();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -156,7 +156,7 @@ public class BootWithMetadataTest
|
|||
assertEquals(toWrite.schema, fromRead.schema);
|
||||
assertEquals(toWrite.directory, fromRead.directory);
|
||||
assertEquals(toWrite.tokenMap, fromRead.tokenMap);
|
||||
assertEquals(toWrite.placements, fromRead.placements);
|
||||
assertEquals(toWrite.placements(), fromRead.placements());
|
||||
assertEquals(toWrite.lockedRanges, fromRead.lockedRanges);
|
||||
assertEquals(toWrite.inProgressSequences, fromRead.inProgressSequences);
|
||||
assertEquals(toWrite.extensions, fromRead.extensions);
|
||||
|
|
|
|||
|
|
@ -348,7 +348,7 @@ public class ClusterMetadataTransformationTest
|
|||
else if (key == TOKEN_MAP)
|
||||
return metadata.tokenMap;
|
||||
else if (key == DATA_PLACEMENTS)
|
||||
return metadata.placements;
|
||||
return metadata.placements();
|
||||
else if (key == LOCKED_RANGES)
|
||||
return metadata.lockedRanges;
|
||||
else if (key == IN_PROGRESS_SEQUENCES)
|
||||
|
|
|
|||
|
|
@ -110,7 +110,7 @@ public class UnregisterTest
|
|||
assertFalse(metadata.directory.allJoinedEndpoints().contains(ep));
|
||||
assertFalse(metadata.directory.allDatacenterRacks().containsKey("dc2"));
|
||||
assertFalse(metadata.directory.knownDatacenters().contains("dc2"));
|
||||
metadata.placements.asMap().forEach((params, placement) -> {
|
||||
metadata.placements().forEach((params, placement) -> {
|
||||
assertFalse(Streams.concat(placement.writes.endpoints.stream(), placement.reads.endpoints.stream()).anyMatch((fr) -> fr.endpoints().contains(ep)));
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -104,7 +104,7 @@ public class GossipHelperTest
|
|||
assertEquals(internal, metadata.directory.addresses.get(nodeId).localAddress);
|
||||
assertEquals(nativeAddress, metadata.directory.addresses.get(nodeId).nativeAddress);
|
||||
|
||||
DataPlacements dp = metadata.placements;
|
||||
DataPlacements dp = metadata.placements();
|
||||
assertEquals(1, dp.get(KSM.params.replication).reads.forToken(token).get().size());
|
||||
assertTrue(dp.get(KSM.params.replication).reads.forToken(token).get().contains(endpoint));
|
||||
assertEquals(1, dp.get(KSM.params.replication).writes.forToken(token).get().size());
|
||||
|
|
@ -196,8 +196,8 @@ public class GossipHelperTest
|
|||
assertEquals(entry.getValue(), metadata.tokenMap.tokens(nodeId).iterator().next());
|
||||
}
|
||||
|
||||
ReplicaGroups reads = metadata.placements.get(KSM_NTS.params.replication).reads;
|
||||
ReplicaGroups writes = metadata.placements.get(KSM_NTS.params.replication).writes;
|
||||
ReplicaGroups reads = metadata.placement(KSM_NTS.params.replication).reads;
|
||||
ReplicaGroups writes = metadata.placement(KSM_NTS.params.replication).writes;
|
||||
assertEquals(reads, writes);
|
||||
// tokens are
|
||||
// dc1: 1: 1000, 3: 3000, 5: 5000, 6: 7000, 7: 9000
|
||||
|
|
|
|||
|
|
@ -109,7 +109,7 @@ public class MetadataSnapshotListenerTest
|
|||
listener.notify(entry, result);
|
||||
ClusterMetadata snapshot = snapshots.getSnapshot(nextEpoch);
|
||||
assertEquals(nextEpoch, snapshot.epoch);
|
||||
assertEquals(toSnapshot.placements, snapshot.placements);
|
||||
assertEquals(toSnapshot.placements(), snapshot.placements());
|
||||
}
|
||||
|
||||
private MetadataSnapshots init()
|
||||
|
|
|
|||
|
|
@ -97,7 +97,7 @@ public class LocalRangesAllSettledTest
|
|||
AllLocalRanges proposed = snapshotAllLocalRanges(LocalRangeStatus.SETTLED, INITIAL_NODES);
|
||||
assertEquals(initial, proposed);
|
||||
// Check against the actual write placements
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, initial, INITIAL_NODES);
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), initial, INITIAL_NODES);
|
||||
|
||||
// Initiate an operation which affects ownership. This will add the MultiStepOperation which encodes any
|
||||
// necessary range movements so subsequent calls to ClusterMetadata::localRangesAllSettled
|
||||
|
|
@ -123,7 +123,7 @@ public class LocalRangesAllSettledTest
|
|||
assertEquals(proposed, finalized);
|
||||
|
||||
// Finally, check against the actual write placements
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, finalized, INITIAL_NODES);
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), finalized, INITIAL_NODES);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
@ -134,7 +134,7 @@ public class LocalRangesAllSettledTest
|
|||
AllLocalRanges proposed = snapshotAllLocalRanges(LocalRangeStatus.SETTLED, INITIAL_NODES);
|
||||
assertEquals(initial, proposed);
|
||||
// Check against the actual write placements
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, initial, INITIAL_NODES);
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), initial, INITIAL_NODES);
|
||||
|
||||
// Initiate an operation which affects ownership. This will add the MultiStepOperation which encodes any
|
||||
// necessary range movements so subsequent calls to ClusterMetadata::localRangesAllSettled
|
||||
|
|
@ -161,7 +161,7 @@ public class LocalRangesAllSettledTest
|
|||
assertEquals(proposed, finalized);
|
||||
|
||||
// Finally, check against the actual write placements
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, finalized, expandedNodes);
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), finalized, expandedNodes);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
@ -172,7 +172,7 @@ public class LocalRangesAllSettledTest
|
|||
AllLocalRanges proposed = snapshotAllLocalRanges(LocalRangeStatus.SETTLED, INITIAL_NODES);
|
||||
assertEquals(initial, proposed);
|
||||
// Check against the actual write placements
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, initial, INITIAL_NODES);
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), initial, INITIAL_NODES);
|
||||
|
||||
// Initiate an operation which affects ownership. This will add the MultiStepOperation which encodes any
|
||||
// necessary range movements so subsequent calls to ClusterMetadata::localRangesAllSettled
|
||||
|
|
@ -201,7 +201,7 @@ public class LocalRangesAllSettledTest
|
|||
assertEquals(proposed, finalized);
|
||||
|
||||
// Finally, check against the actual write placements
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, finalized, INITIAL_NODES);
|
||||
assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), finalized, INITIAL_NODES);
|
||||
}
|
||||
|
||||
private void assertLocalRangesMatchPlacements(DataPlacements placements,
|
||||
|
|
|
|||
|
|
@ -254,6 +254,6 @@ public class OwnershipUtils
|
|||
assert result.isSuccess();
|
||||
workingMetadata = result.success().metadata;
|
||||
}
|
||||
return workingMetadata.placements;
|
||||
return workingMetadata.placements();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -301,7 +301,7 @@ public class InProgressSequenceCancellationTest
|
|||
|
||||
private void assertRelevantMetadata(ClusterMetadata first, ClusterMetadata second)
|
||||
{
|
||||
assertTrue(first.placements.equivalentTo(second.placements));
|
||||
assertTrue(first.placements().equivalentTo(second.placements()));
|
||||
assertTrue(first.directory.equivalentTo(second.directory));
|
||||
assertTrue(first.tokenMap.equivalentTo(second.tokenMap));
|
||||
assertEquals(first.lockedRanges.locked.keySet(), second.lockedRanges.locked.keySet());
|
||||
|
|
|
|||
|
|
@ -113,7 +113,7 @@ public class ProgressBarrierTest extends CMSTestBase
|
|||
// Internally affectedRanges::toPeers uses the same logic as
|
||||
// the progress barrier does to identify the consensus group
|
||||
Set<NodeId> consensusGroup = leave.barrier().affectedRanges.toPeers(ReplicationParams.simple(1),
|
||||
sut.service.metadata().placements,
|
||||
sut.service.metadata().placements(),
|
||||
sut.service.metadata().directory);
|
||||
assertEquals(Set.of(node2.nodeId(), node3.nodeId()), consensusGroup);
|
||||
}
|
||||
|
|
@ -217,7 +217,7 @@ public class ProgressBarrierTest extends CMSTestBase
|
|||
case ALL:
|
||||
{
|
||||
Set<InetAddressAndPort> replicas = metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
|
||||
.stream()
|
||||
.map(n -> metadata.directory.getNodeAddresses(n).broadcastAddress)
|
||||
.collect(Collectors.toSet());
|
||||
|
|
@ -235,7 +235,7 @@ public class ProgressBarrierTest extends CMSTestBase
|
|||
case QUORUM:
|
||||
{
|
||||
Set<InetAddressAndPort> replicas = metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
|
||||
.stream()
|
||||
.map(n -> metadata.directory.getNodeAddresses(n).broadcastAddress)
|
||||
.collect(Collectors.toSet());
|
||||
|
|
@ -253,7 +253,7 @@ public class ProgressBarrierTest extends CMSTestBase
|
|||
case LOCAL_QUORUM:
|
||||
{
|
||||
List<InetAddressAndPort> replicas = new ArrayList<>(metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
|
||||
.stream()
|
||||
.filter((n) -> metadata.directory.location(n).datacenter.equals(dc))
|
||||
.map(n -> metadata.directory.getNodeAddresses(n).broadcastAddress)
|
||||
|
|
@ -277,7 +277,7 @@ public class ProgressBarrierTest extends CMSTestBase
|
|||
{
|
||||
Map<String, Integer> byDc = new HashMap<>();
|
||||
metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
|
||||
.forEach(n -> byDc.compute(metadata.directory.location(n).datacenter,
|
||||
(k, v) -> v == null ? 1 : v + 1));
|
||||
|
||||
|
|
@ -306,7 +306,7 @@ public class ProgressBarrierTest extends CMSTestBase
|
|||
}
|
||||
case ONE:
|
||||
Set<InetAddressAndPort> replicas = metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
|
||||
.toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
|
||||
.stream()
|
||||
.map(n -> metadata.directory.getNodeAddresses(n).broadcastAddress)
|
||||
.collect(Collectors.toSet());
|
||||
|
|
|
|||
|
|
@ -95,8 +95,8 @@ public class EventsMetadataTest
|
|||
|
||||
// should not be in tokenMap (no tokens yet)
|
||||
assertTrue(metadata.tokenMap.tokens(nodeId).isEmpty());
|
||||
assertTrue(metadata.placements.get(KSM.params.replication).writes.byEndpoint().isEmpty());
|
||||
assertTrue(metadata.placements.get(KSM.params.replication).reads.byEndpoint().isEmpty());
|
||||
assertTrue(metadata.placement(KSM.params.replication).writes.byEndpoint().isEmpty());
|
||||
assertTrue(metadata.placement(KSM.params.replication).reads.byEndpoint().isEmpty());
|
||||
assertTrue(metadata.lockedRanges.locked.isEmpty());
|
||||
}
|
||||
|
||||
|
|
@ -130,9 +130,9 @@ public class EventsMetadataTest
|
|||
assertTrue(ClusterMetadata.current().tokenMap.tokens(nodeId).isEmpty());
|
||||
assertEquals(NodeState.BOOTSTRAPPING, ClusterMetadata.current().directory.peerState(nodeId));
|
||||
|
||||
assertTrue(ClusterMetadata.current().placements.get(KSM.params.replication).writes.byEndpoint().containsKey(node1));
|
||||
assertTrue(ClusterMetadata.current().placement(KSM.params.replication).writes.byEndpoint().containsKey(node1));
|
||||
// the first joined node gets added to the read endpoints immediately
|
||||
assertTrue(ClusterMetadata.current().placements.get(KSM.params.replication).reads.byEndpoint().containsKey(node1));
|
||||
assertTrue(ClusterMetadata.current().placement(KSM.params.replication).reads.byEndpoint().containsKey(node1));
|
||||
|
||||
ClusterMetadataService.instance().commit(plan.midJoin);
|
||||
ClusterMetadataService.instance().commit(plan.finishJoin);
|
||||
|
|
@ -152,8 +152,8 @@ public class EventsMetadataTest
|
|||
|
||||
assertTrue(ClusterMetadata.current().tokenMap.tokens(nodeId).isEmpty());
|
||||
assertEquals(NodeState.BOOTSTRAPPING, ClusterMetadata.current().directory.peerState(nodeId));
|
||||
assertTrue(ClusterMetadata.current().placements.get(KSM.params.replication).writes.byEndpoint().containsKey(node2));
|
||||
assertFalse(ClusterMetadata.current().placements.get(KSM.params.replication).reads.byEndpoint().containsKey(node2));
|
||||
assertTrue(ClusterMetadata.current().placement(KSM.params.replication).writes.byEndpoint().containsKey(node2));
|
||||
assertFalse(ClusterMetadata.current().placement(KSM.params.replication).reads.byEndpoint().containsKey(node2));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
@ -178,7 +178,7 @@ public class EventsMetadataTest
|
|||
// no change in metadata after prepareLeave;
|
||||
assertEquals(before.directory, after.directory);
|
||||
assertEquals(before.tokenMap, after.tokenMap);
|
||||
assertEquals(before.placements, after.placements);
|
||||
assertEquals(before.placements(), after.placements());
|
||||
assertEquals(before.schema, after.schema);
|
||||
|
||||
ClusterMetadataService.instance().commit(leave.startLeave);
|
||||
|
|
|
|||
Loading…
Reference in New Issue