From 9297418a14bc6b014b4c23c8ba17b84623ce9347 Mon Sep 17 00:00:00 2001 From: Sam Tunnicliffe Date: Fri, 19 Apr 2024 13:40:14 +0100 Subject: [PATCH] Rename PlacementForRange to ReplicaGroups --- .../cassandra/locator/LocalStrategy.java | 6 +- .../locator/NetworkTopologyStrategy.java | 6 +- .../cassandra/locator/SimpleStrategy.java | 6 +- .../apache/cassandra/tcm/ClusterMetadata.java | 14 ++--- .../tcm/ownership/DataPlacement.java | 52 ++++++++--------- .../tcm/ownership/DataPlacements.java | 10 ++-- ...cementForRange.java => ReplicaGroups.java} | 34 +++++------ .../apache/cassandra/tcm/sequences/Move.java | 4 +- .../tcm/sequences/RemoveNodeStreams.java | 4 +- .../distributed/shared/ClusterUtils.java | 4 +- .../log/MetadataChangeSimulationTest.java | 20 ++++--- .../microbench/PlacementForRangeBench.java | 4 +- .../tcm/compatibility/GossipHelperTest.java | 6 +- .../ownership/UniformRangePlacementTest.java | 58 +++++++++---------- .../InProgressSequenceCancellationTest.java | 10 ++-- 15 files changed, 121 insertions(+), 117 deletions(-) rename src/java/org/apache/cassandra/tcm/ownership/{PlacementForRange.java => ReplicaGroups.java} (94%) diff --git a/src/java/org/apache/cassandra/locator/LocalStrategy.java b/src/java/org/apache/cassandra/locator/LocalStrategy.java index 69193090c4..4032ce1594 100644 --- a/src/java/org/apache/cassandra/locator/LocalStrategy.java +++ b/src/java/org/apache/cassandra/locator/LocalStrategy.java @@ -26,7 +26,7 @@ import org.apache.cassandra.dht.Token; import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tcm.ownership.DataPlacement; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.tcm.ownership.VersionedEndpoints; import org.apache.cassandra.utils.FBUtilities; @@ -65,7 +65,7 @@ public class LocalStrategy extends SystemStrategy { public static final Range entireRange = new Range<>(DatabaseDescriptor.getPartitioner().getMinimumToken(), DatabaseDescriptor.getPartitioner().getMinimumToken()); public static final EndpointsForRange localReplicas = EndpointsForRange.of(new Replica(FBUtilities.getBroadcastAddressAndPort(), entireRange, true)); - public static final DataPlacement placement = new DataPlacement(PlacementForRange.builder().withReplicaGroup(VersionedEndpoints.forRange(Epoch.FIRST, localReplicas)).build(), - PlacementForRange.builder().withReplicaGroup(VersionedEndpoints.forRange(Epoch.FIRST, localReplicas)).build()); + public static final DataPlacement placement = new DataPlacement(ReplicaGroups.builder().withReplicaGroup(VersionedEndpoints.forRange(Epoch.FIRST, localReplicas)).build(), + ReplicaGroups.builder().withReplicaGroup(VersionedEndpoints.forRange(Epoch.FIRST, localReplicas)).build()); } } diff --git a/src/java/org/apache/cassandra/locator/NetworkTopologyStrategy.java b/src/java/org/apache/cassandra/locator/NetworkTopologyStrategy.java index d48ee31610..05bfcfb9ed 100644 --- a/src/java/org/apache/cassandra/locator/NetworkTopologyStrategy.java +++ b/src/java/org/apache/cassandra/locator/NetworkTopologyStrategy.java @@ -49,7 +49,7 @@ import org.apache.cassandra.tcm.membership.Directory; import org.apache.cassandra.tcm.membership.Location; import org.apache.cassandra.tcm.membership.NodeId; import org.apache.cassandra.tcm.ownership.DataPlacement; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.tcm.ownership.TokenMap; import org.apache.cassandra.tcm.ownership.VersionedEndpoints; import org.apache.cassandra.utils.FBUtilities; @@ -194,7 +194,7 @@ public class NetworkTopologyStrategy extends AbstractReplicationStrategy Directory directory, TokenMap tokenMap) { - PlacementForRange.Builder builder = PlacementForRange.builder(); + ReplicaGroups.Builder builder = ReplicaGroups.builder(); for (Range range : ranges) { EndpointsForRange endpointsForRange = calculateNaturalReplicas(range.right, @@ -205,7 +205,7 @@ public class NetworkTopologyStrategy extends AbstractReplicationStrategy builder.withReplicaGroup(VersionedEndpoints.forRange(epoch, endpointsForRange)); } - PlacementForRange built = builder.build(); + ReplicaGroups built = builder.build(); return new DataPlacement(built, built); } diff --git a/src/java/org/apache/cassandra/locator/SimpleStrategy.java b/src/java/org/apache/cassandra/locator/SimpleStrategy.java index 24893e5f1e..380f1a03e5 100644 --- a/src/java/org/apache/cassandra/locator/SimpleStrategy.java +++ b/src/java/org/apache/cassandra/locator/SimpleStrategy.java @@ -41,7 +41,7 @@ import org.apache.cassandra.tcm.compatibility.TokenRingUtils; import org.apache.cassandra.tcm.membership.Directory; import org.apache.cassandra.tcm.membership.NodeId; import org.apache.cassandra.tcm.ownership.DataPlacement; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.tcm.ownership.TokenMap; import org.apache.cassandra.tcm.ownership.VersionedEndpoints; @@ -69,12 +69,12 @@ public class SimpleStrategy extends AbstractReplicationStrategy @Override public DataPlacement calculateDataPlacement(Epoch epoch, List> ranges, ClusterMetadata metadata) { - PlacementForRange.Builder builder = PlacementForRange.builder(); + ReplicaGroups.Builder builder = ReplicaGroups.builder(); for (Range range : ranges) builder.withReplicaGroup(VersionedEndpoints.forRange(epoch, calculateNaturalReplicas(range.right, metadata.tokenMap.tokens(), range, metadata.directory, metadata.tokenMap))); - PlacementForRange built = builder.build(); + ReplicaGroups built = builder.build(); return new DataPlacement(built, built); } diff --git a/src/java/org/apache/cassandra/tcm/ClusterMetadata.java b/src/java/org/apache/cassandra/tcm/ClusterMetadata.java index d01ab11d97..33886bec40 100644 --- a/src/java/org/apache/cassandra/tcm/ClusterMetadata.java +++ b/src/java/org/apache/cassandra/tcm/ClusterMetadata.java @@ -64,7 +64,7 @@ import org.apache.cassandra.tcm.membership.NodeVersion; import org.apache.cassandra.tcm.ownership.DataPlacement; import org.apache.cassandra.tcm.ownership.DataPlacements; import org.apache.cassandra.tcm.ownership.PrimaryRangeComparator; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.tcm.ownership.TokenMap; import org.apache.cassandra.tcm.ownership.VersionedEndpoints; import org.apache.cassandra.tcm.sequences.InProgressSequences; @@ -282,8 +282,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) { - PlacementForRange writes = placements.get(ksm.params.replication).writes; - PlacementForRange reads = placements.get(ksm.params.replication).reads; + ReplicaGroups writes = placements.get(ksm.params.replication).writes; + ReplicaGroups reads = placements.get(ksm.params.replication).reads; if (ksm.params.replication.isMeta()) return !reads.equals(writes); return !reads.forToken(token).equals(writes.forToken(token)); @@ -292,8 +292,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) { - PlacementForRange writes = placements.get(ksm.params.replication).writes; - PlacementForRange reads = placements.get(ksm.params.replication).reads; + ReplicaGroups writes = placements.get(ksm.params.replication).writes; + ReplicaGroups reads = placements.get(ksm.params.replication).reads; return !writes.byEndpoint().get(endpoint).equals(reads.byEndpoint().get(endpoint)); } @@ -311,8 +311,8 @@ public class ClusterMetadata public Map, VersionedEndpoints.ForRange> pendingRanges(KeyspaceMetadata metadata) { Map, VersionedEndpoints.ForRange> map = new HashMap<>(); - PlacementForRange writes = placements.get(metadata.params.replication).writes; - PlacementForRange reads = placements.get(metadata.params.replication).reads; + ReplicaGroups writes = placements.get(metadata.params.replication).writes; + ReplicaGroups reads = placements.get(metadata.params.replication).reads; // first, pending ranges as the result of range splitting or merging // i.e. new ranges being created through join/leave diff --git a/src/java/org/apache/cassandra/tcm/ownership/DataPlacement.java b/src/java/org/apache/cassandra/tcm/ownership/DataPlacement.java index 28d5d5b1ad..f42d8b5479 100644 --- a/src/java/org/apache/cassandra/tcm/ownership/DataPlacement.java +++ b/src/java/org/apache/cassandra/tcm/ownership/DataPlacement.java @@ -46,16 +46,16 @@ public class DataPlacement return replication.isMeta() ? metaKeyspaceSerializer : globalSerializer; } - private static final DataPlacement EMPTY = new DataPlacement(PlacementForRange.EMPTY, PlacementForRange.EMPTY); + private static final DataPlacement EMPTY = new DataPlacement(ReplicaGroups.EMPTY, ReplicaGroups.EMPTY); // TODO make tree of just EndpointsForRange, navigable by EFR.range() // TODO combine peers into a single entity with one vote in any quorum // (e.g. old & new peer must both respond to count one replica) - public final PlacementForRange reads; - public final PlacementForRange writes; + public final ReplicaGroups reads; + public final ReplicaGroups writes; - public DataPlacement(PlacementForRange reads, - PlacementForRange writes) + public DataPlacement(ReplicaGroups reads, + ReplicaGroups writes) { this.reads = reads; this.writes = writes; @@ -76,14 +76,14 @@ public class DataPlacement public DataPlacement combineReplicaGroups(DataPlacement other) { - return new DataPlacement(PlacementForRange.builder() - .withReplicaGroups(reads.endpoints) - .withReplicaGroups(other.reads.endpoints) - .build(), - PlacementForRange.builder() - .withReplicaGroups(writes.endpoints) - .withReplicaGroups(other.writes.endpoints) - .build()); + return new DataPlacement(ReplicaGroups.builder() + .withReplicaGroups(reads.endpoints) + .withReplicaGroups(other.reads.endpoints) + .build(), + ReplicaGroups.builder() + .withReplicaGroups(writes.endpoints) + .withReplicaGroups(other.writes.endpoints) + .build()); } public PlacementDeltas.PlacementDelta difference(DataPlacement next) @@ -94,8 +94,8 @@ public class DataPlacement public DataPlacement splitRangesForPlacement(List tokens) { - return new DataPlacement(PlacementForRange.splitRangesForPlacement(tokens, reads), - PlacementForRange.splitRangesForPlacement(tokens, writes)); + return new DataPlacement(ReplicaGroups.splitRangesForPlacement(tokens, reads), + ReplicaGroups.splitRangesForPlacement(tokens, writes)); } public static DataPlacement empty() @@ -105,8 +105,8 @@ public class DataPlacement public static Builder builder() { - return new Builder(PlacementForRange.builder(), - PlacementForRange.builder()); + return new Builder(ReplicaGroups.builder(), + ReplicaGroups.builder()); } public Builder unbuild() @@ -115,10 +115,10 @@ public class DataPlacement } public static class Builder { - public final PlacementForRange.Builder reads; - public final PlacementForRange.Builder writes; + public final ReplicaGroups.Builder reads; + public final ReplicaGroups.Builder writes; - public Builder(PlacementForRange.Builder reads, PlacementForRange.Builder writes) + public Builder(ReplicaGroups.Builder reads, ReplicaGroups.Builder writes) { this.reads = reads; this.writes = writes; @@ -194,21 +194,21 @@ public class DataPlacement public void serialize(DataPlacement t, DataOutputPlus out, Version version) throws IOException { - PlacementForRange.serializer.serialize(t.reads, out, partitioner, version); - PlacementForRange.serializer.serialize(t.writes, out, partitioner, version); + ReplicaGroups.serializer.serialize(t.reads, out, partitioner, version); + ReplicaGroups.serializer.serialize(t.writes, out, partitioner, version); } public DataPlacement deserialize(DataInputPlus in, Version version) throws IOException { - PlacementForRange reads = PlacementForRange.serializer.deserialize(in, partitioner, version); - PlacementForRange writes = PlacementForRange.serializer.deserialize(in, partitioner, version); + ReplicaGroups reads = ReplicaGroups.serializer.deserialize(in, partitioner, version); + ReplicaGroups writes = ReplicaGroups.serializer.deserialize(in, partitioner, version); return new DataPlacement(reads, writes); } public long serializedSize(DataPlacement t, Version version) { - return PlacementForRange.serializer.serializedSize(t.reads, partitioner, version) + - PlacementForRange.serializer.serializedSize(t.writes, partitioner, version); + return ReplicaGroups.serializer.serializedSize(t.reads, partitioner, version) + + ReplicaGroups.serializer.serializedSize(t.writes, partitioner, version); } } } diff --git a/src/java/org/apache/cassandra/tcm/ownership/DataPlacements.java b/src/java/org/apache/cassandra/tcm/ownership/DataPlacements.java index f737e3dbbf..988d2b1bcb 100644 --- a/src/java/org/apache/cassandra/tcm/ownership/DataPlacements.java +++ b/src/java/org/apache/cassandra/tcm/ownership/DataPlacements.java @@ -101,9 +101,9 @@ public class DataPlacements extends ReplicationMap implements Met DatabaseDescriptor.getPartitioner().getMinimumToken(), DatabaseDescriptor.getPartitioner().getMinimumToken(), true)); - PlacementForRange placement = PlacementForRange.builder(1) - .withReplicaGroup(VersionedEndpoints.forRange(Epoch.EMPTY, endpoints)) - .build(); + ReplicaGroups placement = ReplicaGroups.builder(1) + .withReplicaGroup(VersionedEndpoints.forRange(Epoch.EMPTY, endpoints)) + .build(); LOCAL_PLACEMENT = new DataPlacement(placement, placement); } return LOCAL_PLACEMENT; @@ -149,12 +149,12 @@ public class DataPlacements extends ReplicationMap implements Met builder.with(params, placement); else { - PlacementForRange.Builder reads = PlacementForRange.builder(placement.reads.size()); + ReplicaGroups.Builder reads = ReplicaGroups.builder(placement.reads.size()); placement.reads.endpoints.forEach((endpoints) -> { reads.withReplicaGroup(VersionedEndpoints.forRange(endpoints.lastModified(), endpoints.get().sorted(comparator))); }); - PlacementForRange.Builder writes = PlacementForRange.builder(placement.writes.size()); + ReplicaGroups.Builder writes = ReplicaGroups.builder(placement.writes.size()); placement.writes.endpoints.forEach((endpoints) -> { writes.withReplicaGroup(VersionedEndpoints.forRange(endpoints.lastModified(), endpoints.get().sorted(comparator))); diff --git a/src/java/org/apache/cassandra/tcm/ownership/PlacementForRange.java b/src/java/org/apache/cassandra/tcm/ownership/ReplicaGroups.java similarity index 94% rename from src/java/org/apache/cassandra/tcm/ownership/PlacementForRange.java rename to src/java/org/apache/cassandra/tcm/ownership/ReplicaGroups.java index 4023e49b69..66a3dc8a21 100644 --- a/src/java/org/apache/cassandra/tcm/ownership/PlacementForRange.java +++ b/src/java/org/apache/cassandra/tcm/ownership/ReplicaGroups.java @@ -48,7 +48,7 @@ import org.apache.cassandra.utils.AsymmetricOrdering; import static org.apache.cassandra.db.TypeSizes.sizeof; -public class PlacementForRange +public class ReplicaGroups { private static final AsymmetricOrdering, Token> ordering = new AsymmetricOrdering<>() { @@ -72,12 +72,12 @@ public class PlacementForRange }; public static final Serializer serializer = new Serializer(); - public static final PlacementForRange EMPTY = PlacementForRange.builder().build(); + public static final ReplicaGroups EMPTY = ReplicaGroups.builder().build(); public final ImmutableList> ranges; public final ImmutableList endpoints; - public PlacementForRange(Map, VersionedEndpoints.ForRange> replicaGroups) + public ReplicaGroups(Map, VersionedEndpoints.ForRange> replicaGroups) { ImmutableList.Builder> rangesBuilder = ImmutableList.builderWithExpectedSize(replicaGroups.size()); ImmutableList.Builder endpointsBuilder = ImmutableList.builderWithExpectedSize(replicaGroups.size()); @@ -146,7 +146,7 @@ public class PlacementForRange return forRange(token).forToken(token); } - public Delta difference(PlacementForRange next) + public Delta difference(ReplicaGroups next) { RangesByEndpoint oldMap = this.byEndpoint(); RangesByEndpoint newMap = next.byEndpoint(); @@ -179,7 +179,7 @@ public class PlacementForRange return builder.build(); } - public PlacementForRange withCappedLastModified(Epoch lastModified) + public ReplicaGroups withCappedLastModified(Epoch lastModified) { SortedMap, VersionedEndpoints.ForRange> copy = new TreeMap<>(); for (int i = 0; i < ranges.size(); i++) @@ -190,7 +190,7 @@ public class PlacementForRange forRange = forRange.withLastModified(lastModified); copy.put(range, forRange); } - return new PlacementForRange(copy); + return new ReplicaGroups(copy); } @@ -262,12 +262,12 @@ public class PlacementForRange } @VisibleForTesting - public static PlacementForRange splitRangesForPlacement(List tokens, PlacementForRange placement) + public static ReplicaGroups splitRangesForPlacement(List tokens, ReplicaGroups placement) { if (placement.ranges.isEmpty()) return placement; - Builder newPlacement = PlacementForRange.builder(); + Builder newPlacement = ReplicaGroups.builder(); List eprs = new ArrayList<>(placement.endpoints); eprs.sort(Comparator.comparing(a -> a.range().left)); Token min = eprs.get(0).range().left; @@ -426,15 +426,15 @@ public class PlacementForRange return this; } - public PlacementForRange build() + public ReplicaGroups build() { - return new PlacementForRange(this.replicaGroups); + return new ReplicaGroups(this.replicaGroups); } } - public static class Serializer implements PartitionerAwareMetadataSerializer + public static class Serializer implements PartitionerAwareMetadataSerializer { - public void serialize(PlacementForRange t, DataOutputPlus out, IPartitioner partitioner, Version version) throws IOException + public void serialize(ReplicaGroups t, DataOutputPlus out, IPartitioner partitioner, Version version) throws IOException { out.writeInt(t.ranges.size()); @@ -458,7 +458,7 @@ public class PlacementForRange } } - public PlacementForRange deserialize(DataInputPlus in, IPartitioner partitioner, Version version) throws IOException + public ReplicaGroups deserialize(DataInputPlus in, IPartitioner partitioner, Version version) throws IOException { int groupCount = in.readInt(); Map, VersionedEndpoints.ForRange> result = Maps.newHashMapWithExpectedSize(groupCount); @@ -492,10 +492,10 @@ public class PlacementForRange EndpointsForRange efr = EndpointsForRange.copyOf(replicas); result.put(range, VersionedEndpoints.forRange(lastModified, efr)); } - return new PlacementForRange(result); + return new ReplicaGroups(result); } - public long serializedSize(PlacementForRange t, IPartitioner partitioner, Version version) + public long serializedSize(ReplicaGroups t, IPartitioner partitioner, Version version) { long size = sizeof(t.ranges.size()); for (int i = 0; i < t.ranges.size(); i++) @@ -525,8 +525,8 @@ public class PlacementForRange public boolean equals(Object o) { if (this == o) return true; - if (!(o instanceof PlacementForRange)) return false; - PlacementForRange that = (PlacementForRange) o; + if (!(o instanceof ReplicaGroups)) return false; + ReplicaGroups that = (ReplicaGroups) o; return Objects.equals(ranges, that.ranges) && Objects.equals(endpoints, that.endpoints); } diff --git a/src/java/org/apache/cassandra/tcm/sequences/Move.java b/src/java/org/apache/cassandra/tcm/sequences/Move.java index 898cfabc3b..4b28d7b8aa 100644 --- a/src/java/org/apache/cassandra/tcm/sequences/Move.java +++ b/src/java/org/apache/cassandra/tcm/sequences/Move.java @@ -61,7 +61,7 @@ import org.apache.cassandra.tcm.membership.NodeState; import org.apache.cassandra.tcm.ownership.DataPlacements; import org.apache.cassandra.tcm.ownership.MovementMap; import org.apache.cassandra.tcm.ownership.PlacementDeltas; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.tcm.serialization.AsymmetricMetadataSerializer; import org.apache.cassandra.tcm.serialization.MetadataSerializer; import org.apache.cassandra.tcm.serialization.Version; @@ -334,7 +334,7 @@ public class Move extends MultiStepOperation MovementMap.Builder allMovements = MovementMap.builder(); toStart.forEach((params, delta) -> { RangesByEndpoint targets = delta.writes.additions; - PlacementForRange oldOwners = placements.get(params).reads; + ReplicaGroups oldOwners = placements.get(params).reads; EndpointsByReplica.Builder movements = new EndpointsByReplica.Builder(); targets.flattenValues().forEach(destination -> { SourceHolder sources = new SourceHolder(fd, destination, toSplitRanges.get(params), strictConsistency); diff --git a/src/java/org/apache/cassandra/tcm/sequences/RemoveNodeStreams.java b/src/java/org/apache/cassandra/tcm/sequences/RemoveNodeStreams.java index b5951b6183..d6e8ca3127 100644 --- a/src/java/org/apache/cassandra/tcm/sequences/RemoveNodeStreams.java +++ b/src/java/org/apache/cassandra/tcm/sequences/RemoveNodeStreams.java @@ -39,7 +39,7 @@ import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.membership.NodeId; import org.apache.cassandra.tcm.ownership.MovementMap; import org.apache.cassandra.tcm.ownership.PlacementDeltas; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import static org.apache.cassandra.streaming.StreamOperation.RESTORE_REPLICA_COUNT; @@ -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 - PlacementForRange currentPlacements = metadata.placements.get(params).reads; + ReplicaGroups currentPlacements = metadata.placements.get(params).reads; startWriteAdditions.flattenValues() .forEach(newReplica -> { EndpointsForRange.Builder candidateBuilder = new EndpointsForRange.Builder(newReplica.range()); diff --git a/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java b/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java index b160b2600f..f6628d3047 100644 --- a/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java +++ b/test/distributed/org/apache/cassandra/distributed/shared/ClusterUtils.java @@ -78,7 +78,7 @@ import org.apache.cassandra.tcm.Commit; import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tcm.Transformation; import org.apache.cassandra.tcm.membership.NodeId; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.utils.Isolated; import org.apache.cassandra.utils.concurrent.AsyncPromise; import org.apache.cassandra.utils.concurrent.CountDownLatch; @@ -388,7 +388,7 @@ public class ClusterUtils StringBuilder builder = new StringBuilder(); builder.append("'keyspace' { 'name':").append(keyspace.name).append("', "); builder.append("'reads':['"); - PlacementForRange placement = metadata.placements.get(keyspace.params.replication).reads; + ReplicaGroups placement = metadata.placements.get(keyspace.params.replication).reads; builder.append(byEndpoint ? placement.toStringByEndpoint() : placement.toString()); builder.append("'], 'writes':['"); placement = metadata.placements.get(keyspace.params.replication).writes; diff --git a/test/distributed/org/apache/cassandra/distributed/test/log/MetadataChangeSimulationTest.java b/test/distributed/org/apache/cassandra/distributed/test/log/MetadataChangeSimulationTest.java index 33b6ce687e..26d1cff7ac 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/log/MetadataChangeSimulationTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/log/MetadataChangeSimulationTest.java @@ -65,7 +65,7 @@ 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.ownership.DataPlacements; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.tcm.ownership.VersionedEndpoints; import org.apache.cassandra.tcm.transformations.Register; import org.apache.cassandra.tcm.transformations.TriggerSnapshot; @@ -700,7 +700,7 @@ public class MetadataChangeSimulationTest extends CMSTestBase return sb.toString(); } - public static void match(PlacementForRange actual, Map> predicted) throws Throwable + public static void match(ReplicaGroups actual, Map> predicted) throws Throwable { Map, VersionedEndpoints.ForRange> actualGroups = actual.asMap(); assert predicted.size() == actualGroups.size() : @@ -825,24 +825,28 @@ public class MetadataChangeSimulationTest extends CMSTestBase validatePlacementsInternal(rf, modelState.inFlightOperations, expectedRanges, actualPlacements.writes, true); } - public static void validateTransientStatus(PlacementForRange reads, PlacementForRange writes) + public static void validateTransientStatus(ReplicaGroups reads, ReplicaGroups writes) { // No node should ever be a FULL read replica but a TRANSIENT write replica for the same range Map, List> invalid = new HashMap<>(); - reads.replicaGroups().forEach((range, readGroup) -> { + for (int i = 0; i < reads.ranges.size(); i++) + { + Range range = reads.ranges.get(i); + VersionedEndpoints.ForRange readGroup = reads.endpoints.get(i); Map writeGroup = writes.forRange(range).get().byEndpoint(); + readGroup.forEach(r -> { if (r.isFull()) { Replica w = writeGroup.get(r.endpoint()); if (w != null && w.isTransient()) { - List i = invalid.computeIfAbsent(range, ignore -> new ArrayList<>()); - i.add(w); + List replicas = invalid.computeIfAbsent(range, ignore -> new ArrayList<>()); + replicas.add(w); } } }); - }); + } assertTrue(() -> String.format("Found replicas with invalid transient/full status within a given range. " + "The following were found with the same instance having TRANSIENT status for writes, but " + "FULL status for reads, which can cause consistency violations. %n%s", invalid), @@ -854,7 +858,7 @@ public class MetadataChangeSimulationTest extends CMSTestBase Assert.assertEquals(new TreeSet<>(l), new TreeSet<>(r)); } - public static void validatePlacementsInternal(ReplicationFactor rf, List opStates, List> expectedRanges, PlacementForRange placements, boolean allowPending) + public static void validatePlacementsInternal(ReplicationFactor rf, List opStates, List> expectedRanges, ReplicaGroups placements, boolean allowPending) { int overreplicated = 0; for (Range range : expectedRanges) diff --git a/test/microbench/org/apache/cassandra/test/microbench/PlacementForRangeBench.java b/test/microbench/org/apache/cassandra/test/microbench/PlacementForRangeBench.java index da845c5f90..407559509c 100644 --- a/test/microbench/org/apache/cassandra/test/microbench/PlacementForRangeBench.java +++ b/test/microbench/org/apache/cassandra/test/microbench/PlacementForRangeBench.java @@ -39,7 +39,7 @@ import org.apache.cassandra.tcm.membership.Location; import org.apache.cassandra.tcm.membership.NodeAddresses; import org.apache.cassandra.tcm.membership.NodeId; import org.apache.cassandra.tcm.membership.NodeVersion; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.tcm.ownership.PlacementProvider; import org.apache.cassandra.tcm.ownership.UniformRangePlacement; import org.apache.cassandra.tcm.transformations.UnsafeJoin; @@ -69,7 +69,7 @@ public class PlacementForRangeBench { static Token [] queryTokens = new Token[5000]; static Random random = new Random(1); - static PlacementForRange pfr; + static ReplicaGroups pfr; /* new: PlacementForRangeBench.bench avgt 5 0,317 ± 0,037 ms/op diff --git a/test/unit/org/apache/cassandra/tcm/compatibility/GossipHelperTest.java b/test/unit/org/apache/cassandra/tcm/compatibility/GossipHelperTest.java index 329770db07..3ff7c03e83 100644 --- a/test/unit/org/apache/cassandra/tcm/compatibility/GossipHelperTest.java +++ b/test/unit/org/apache/cassandra/tcm/compatibility/GossipHelperTest.java @@ -46,7 +46,7 @@ import org.apache.cassandra.schema.Keyspaces; import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.ownership.DataPlacements; import org.apache.cassandra.tcm.membership.NodeId; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.utils.CassandraVersion; import static org.apache.cassandra.gms.ApplicationState.*; @@ -157,8 +157,8 @@ public class GossipHelperTest assertEquals(entry.getValue(), metadata.tokenMap.tokens(nodeId).iterator().next()); } - PlacementForRange reads = metadata.placements.get(KSM_NTS.params.replication).reads; - PlacementForRange writes = metadata.placements.get(KSM_NTS.params.replication).writes; + ReplicaGroups reads = metadata.placements.get(KSM_NTS.params.replication).reads; + ReplicaGroups writes = metadata.placements.get(KSM_NTS.params.replication).writes; assertEquals(reads, writes); // tokens are // dc1: 1: 1000, 3: 3000, 5: 5000, 6: 7000, 7: 9000 diff --git a/test/unit/org/apache/cassandra/tcm/ownership/UniformRangePlacementTest.java b/test/unit/org/apache/cassandra/tcm/ownership/UniformRangePlacementTest.java index 6f54491e3f..c987eff2e8 100644 --- a/test/unit/org/apache/cassandra/tcm/ownership/UniformRangePlacementTest.java +++ b/test/unit/org/apache/cassandra/tcm/ownership/UniformRangePlacementTest.java @@ -50,12 +50,12 @@ public class UniformRangePlacementTest // which is either the min or max value in the token space. Any other single token would produce two ranges - // (MIN, t] & (t, MAX] - but because (x, x] denotes a wraparound, this case produces (MIN, MAX] and we need to // verify that we can safely split that when more tokens are introduced. This test supposes MIN = 0, MAX = 100 - PlacementForRange before = PlacementForRange.builder() - .withReplicaGroup(VersionedEndpoints.forRange(Epoch.EMPTY, rg(0, 100, 1, 2, 3))) - .build(); + ReplicaGroups before = ReplicaGroups.builder() + .withReplicaGroup(VersionedEndpoints.forRange(Epoch.EMPTY, rg(0, 100, 1, 2, 3))) + .build(); // existing token is MIN (i.e. 0 for the purposes of this test) List tokens = ImmutableList.of(token(0), token(30), token(60), token(90)); - PlacementForRange after = PlacementForRange.splitRangesForPlacement(tokens, before); + ReplicaGroups after = ReplicaGroups.splitRangesForPlacement(tokens, before); assertPlacement(after, rg(0, 30, 1, 2, 3), rg(30, 60, 1, 2, 3), @@ -65,7 +65,7 @@ public class UniformRangePlacementTest // existing token is MAX (i.e. 100 for the purposes of this test). tokens = ImmutableList.of(token(30), token(60), token(90), token(100)); - after = PlacementForRange.splitRangesForPlacement(tokens, before); + after = ReplicaGroups.splitRangesForPlacement(tokens, before); assertPlacement(after, rg(0, 30, 1, 2, 3), rg(30, 60, 1, 2, 3), @@ -80,10 +80,10 @@ public class UniformRangePlacementTest // (100,200] : (n1,n2,n3); // (200,300] : (n1,n2,n3); // (300,400] : (n1,n2,n3); - PlacementForRange before = initialPlacement(); + ReplicaGroups before = initialPlacement(); // split (100,200] to (100,150], (150,200] List tokens = ImmutableList.of(token(100), token(150), token(200), token(300)); - PlacementForRange after = PlacementForRange.splitRangesForPlacement(tokens, before); + ReplicaGroups after = ReplicaGroups.splitRangesForPlacement(tokens, before); assertPlacement(after, rg(0, 100, 1, 2, 3), rg(100, 150, 1, 2, 3), @@ -99,11 +99,11 @@ public class UniformRangePlacementTest // (100,200] : (n1,n2,n3); // (200,300] : (n1,n2,n3); // (300,400] : (n1,n2,n3); - PlacementForRange before = initialPlacement(); + ReplicaGroups before = initialPlacement(); // split (100,200] to (100,150],(150,200] // and (200,300] to (200,250],(250,300] List tokens = ImmutableList.of(token(100), token(150), token(200), token(250), token(300)); - PlacementForRange after = PlacementForRange.splitRangesForPlacement(tokens, before); + ReplicaGroups after = ReplicaGroups.splitRangesForPlacement(tokens, before); assertPlacement(after, rg(0, 100, 1, 2, 3), rg(100, 150, 1, 2, 3), @@ -120,10 +120,10 @@ public class UniformRangePlacementTest // (100,200] : (n1,n2,n3); // (200,300] : (n1,n2,n3); // (300,400] : (n1,n2,n3); - PlacementForRange before = initialPlacement(); + ReplicaGroups before = initialPlacement(); // split (100,200] to (100,125],(125,150],(150,200] List tokens = ImmutableList.of(token(100), token(125), token(150), token(200), token(300)); - PlacementForRange after = PlacementForRange.splitRangesForPlacement(tokens, before); + ReplicaGroups after = ReplicaGroups.splitRangesForPlacement(tokens, before); assertPlacement(after, rg(0, 100, 1, 2, 3), rg(100, 125, 1, 2, 3), @@ -136,11 +136,11 @@ public class UniformRangePlacementTest @Test public void testSplitMultipleRangesMultipleTimes() { - PlacementForRange before = initialPlacement(); + ReplicaGroups before = initialPlacement(); // split (100,200] to (100,125],(125,150],(150,200] // and (200,300] to (200,225],(225,250],(250,300] List tokens = ImmutableList.of(token(100), token(125), token(150), token(200), token(225), token(250), token(300)); - PlacementForRange after = PlacementForRange.splitRangesForPlacement(tokens, before); + ReplicaGroups after = ReplicaGroups.splitRangesForPlacement(tokens, before); assertPlacement(after, rg(0, 100, 1, 2, 3), rg(100, 125, 1, 2, 3), @@ -159,10 +159,10 @@ public class UniformRangePlacementTest // (100,200] : (n1,n2,n3); // (200,300] : (n1,n2,n3); // (300,400] : (n1,n2,n3); - PlacementForRange before = initialPlacement(); + ReplicaGroups before = initialPlacement(); // split (300,400] to (300,325],(325,350],(350,400] List tokens = ImmutableList.of(token(100), token(200), token(300), token(325), token(350)); - PlacementForRange after = PlacementForRange.splitRangesForPlacement(tokens, before); + ReplicaGroups after = ReplicaGroups.splitRangesForPlacement(tokens, before); assertPlacement(after, rg(0, 100, 1, 2, 3), rg(100, 200, 1, 2, 3), @@ -179,10 +179,10 @@ public class UniformRangePlacementTest // (100,200] : (n1,n2,n3); // (200,300] : (n1,n2,n3); // (300,400] : (n1,n2,n3); - PlacementForRange before = initialPlacement(); + ReplicaGroups before = initialPlacement(); // split (0,100] to (0,25],(25,50],(50,100] List tokens = ImmutableList.of(token(25), token(50), token(100), token(200), token(300)); - PlacementForRange after = PlacementForRange.splitRangesForPlacement(tokens, before); + ReplicaGroups after = ReplicaGroups.splitRangesForPlacement(tokens, before); assertPlacement(after, rg(0, 25, 1, 2, 3), rg(25, 50, 1, 2, 3), @@ -199,22 +199,22 @@ public class UniformRangePlacementTest rg(100, 200, 1, 2, 3), rg(200, 300, 1, 2, 3), rg(300, 400, 1, 2, 3) }; - PlacementForRange p1 = PlacementForRange.builder() - .withReplicaGroups(Arrays.asList(firstGroups).stream().map(this::v).collect(Collectors.toList())) - .build(); + ReplicaGroups p1 = ReplicaGroups.builder() + .withReplicaGroups(Arrays.asList(firstGroups).stream().map(this::v).collect(Collectors.toList())) + .build(); EndpointsForRange[] secondGroups = { rg(0, 100, 2, 3, 4), rg(100, 200, 2, 3, 5), rg(200, 300, 2, 3, 6), rg(300, 400, 2, 3, 7) }; - PlacementForRange p2 = PlacementForRange.builder() - .withReplicaGroups(Arrays.asList(secondGroups).stream().map(this::v).collect(Collectors.toList())) - .build(); + ReplicaGroups p2 = ReplicaGroups.builder() + .withReplicaGroups(Arrays.asList(secondGroups).stream().map(this::v).collect(Collectors.toList())) + .build(); ReplicationParams params = ReplicationParams.simple(1); DataPlacements map1 = DataPlacements.builder(1).with(params, new DataPlacement(p1, p1)).build(); DataPlacements map2 = DataPlacements.builder(1).with(params, new DataPlacement(p2, p2)).build(); DataPlacement p3 = map1.combineReplicaGroups(map2).get(params); - for (PlacementForRange placement : new PlacementForRange[]{ p3.reads, p3.writes }) + for (ReplicaGroups placement : new ReplicaGroups[]{ p3.reads, p3.writes }) { assertPlacement(placement, rg(0, 100, 1, 2, 3, 4), @@ -269,20 +269,20 @@ public class UniformRangePlacementTest assertEquals(2, newPlacement.writes.size()); } - private PlacementForRange initialPlacement() + private ReplicaGroups initialPlacement() { EndpointsForRange[] initialGroups = { rg(0, 100, 1, 2, 3), rg(100, 200, 1, 2, 3), rg(200, 300, 1, 2, 3), rg(300, 400, 1, 2, 3) }; - PlacementForRange placement = PlacementForRange.builder() - .withReplicaGroups(Arrays.stream(initialGroups).map(this::v).collect(Collectors.toList())) - .build(); + ReplicaGroups placement = ReplicaGroups.builder() + .withReplicaGroups(Arrays.stream(initialGroups).map(this::v).collect(Collectors.toList())) + .build(); assertPlacement(placement, initialGroups); return placement; } - private void assertPlacement(PlacementForRange placement, EndpointsForRange...expected) + private void assertPlacement(ReplicaGroups placement, EndpointsForRange...expected) { Collection replicaGroups = placement.endpoints.stream().map(VersionedEndpoints.ForRange::get).collect(Collectors.toList()); assertEquals(replicaGroups.size(), expected.length); diff --git a/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java b/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java index 80a806863f..1c1a44296b 100644 --- a/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java +++ b/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java @@ -47,7 +47,7 @@ import org.apache.cassandra.tcm.membership.NodeVersion; import org.apache.cassandra.tcm.ownership.DataPlacement; import org.apache.cassandra.tcm.ownership.DataPlacements; import org.apache.cassandra.tcm.ownership.PlacementDeltas; -import org.apache.cassandra.tcm.ownership.PlacementForRange; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; import org.apache.cassandra.tcm.transformations.PrepareJoin; import org.apache.cassandra.tcm.transformations.PrepareLeave; import org.apache.cassandra.tcm.transformations.PrepareReplace; @@ -321,8 +321,8 @@ public class InProgressSequenceCancellationTest first.asMap().forEach((params, placement) -> { DataPlacement otherPlacement = second.get(params); - PlacementForRange r1 = placement.reads; - PlacementForRange r2 = otherPlacement.reads; + ReplicaGroups r1 = placement.reads; + ReplicaGroups r2 = otherPlacement.reads; assertEquals(r1.ranges, r2.ranges); r1.forEach((range, e1) -> { EndpointsForRange e2 = r2.forRange(range).get(); @@ -330,8 +330,8 @@ public class InProgressSequenceCancellationTest assertTrue(e1.get().stream().allMatch(e2::contains)); }); - PlacementForRange w1 = placement.reads; - PlacementForRange w2 = otherPlacement.reads; + ReplicaGroups w1 = placement.reads; + ReplicaGroups w2 = otherPlacement.reads; assertEquals(w1.ranges, w2.ranges); w1.forEach((range, e1) -> { EndpointsForRange e2 = w2.forRange(range).get();