Combine PlacementsForRange and ReplicaGroups

This commit is contained in:
Sam Tunnicliffe 2024-04-19 13:30:35 +01:00 committed by Marcus Eriksson
parent 99147a0f06
commit c958788ed8
11 changed files with 112 additions and 152 deletions

View File

@ -323,7 +323,7 @@ public class ClusterMetadata
// next, ranges where the ranges themselves are not changing, but the replicas are
// i.e. replacement or RF increase
writes.replicaGroups().forEach((range, endpoints) -> {
writes.forEach((range, endpoints) -> {
VersionedEndpoints.ForRange readGroup = reads.forRange(range);
if (!readGroup.equals(endpoints))
map.put(range, VersionedEndpoints.forRange(endpoints.lastModified(),

View File

@ -77,12 +77,12 @@ public class DataPlacement
public DataPlacement combineReplicaGroups(DataPlacement other)
{
return new DataPlacement(PlacementForRange.builder()
.withReplicaGroups(reads.replicaGroups().endpoints)
.withReplicaGroups(other.reads.replicaGroups().endpoints)
.withReplicaGroups(reads.endpoints)
.withReplicaGroups(other.reads.endpoints)
.build(),
PlacementForRange.builder()
.withReplicaGroups(writes.replicaGroups().endpoints)
.withReplicaGroups(other.writes.replicaGroups().endpoints)
.withReplicaGroups(writes.endpoints)
.withReplicaGroups(other.writes.endpoints)
.build());
}

View File

@ -149,13 +149,13 @@ public class DataPlacements extends ReplicationMap<DataPlacement> implements Met
builder.with(params, placement);
else
{
PlacementForRange.Builder reads = PlacementForRange.builder(placement.reads.replicaGroups().size());
placement.reads.replicaGroups().endpoints.forEach((endpoints) -> {
PlacementForRange.Builder reads = PlacementForRange.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.replicaGroups().size());
placement.writes.replicaGroups().endpoints.forEach((endpoints) -> {
PlacementForRange.Builder writes = PlacementForRange.builder(placement.writes.size());
placement.writes.endpoints.forEach((endpoints) -> {
writes.withReplicaGroup(VersionedEndpoints.forRange(endpoints.lastModified(),
endpoints.get().sorted(comparator)));
});

View File

@ -50,11 +50,32 @@ import static org.apache.cassandra.db.TypeSizes.sizeof;
public class PlacementForRange
{
public static final Serializer serializer = new Serializer();
private static final AsymmetricOrdering<Range<Token>, Token> ordering = new AsymmetricOrdering<>()
{
@Override
public int compare(Range<Token> left, Range<Token> right)
{
return left.compareTo(right);
}
@Override
public int compareAsymmetric(Range<Token> range, Token token)
{
if (token.isMinimum() && !range.right.isMinimum())
return -1;
if (range.left.compareTo(token) >= 0)
return 1;
if (!range.right.isMinimum() && range.right.compareTo(token) < 0)
return -1;
return 0;
}
};
public static final Serializer serializer = new Serializer();
public static final PlacementForRange EMPTY = PlacementForRange.builder().build();
private final ReplicaGroups replicaGroups;
public final ImmutableList<Range<Token>> ranges;
public final ImmutableList<VersionedEndpoints.ForRange> endpoints;
public PlacementForRange(Map<Range<Token>, VersionedEndpoints.ForRange> replicaGroups)
{
@ -69,19 +90,14 @@ public class PlacementForRange
rangesBuilder.add(entry.getKey());
endpointsBuilder.add(entry.getValue());
}
this.replicaGroups = new ReplicaGroups(rangesBuilder.build(), endpointsBuilder.build());
}
@VisibleForTesting
public ReplicaGroups replicaGroups()
{
return replicaGroups;
this.ranges = rangesBuilder.build();
this.endpoints = endpointsBuilder.build();
}
@VisibleForTesting
public List<Range<Token>> ranges()
{
List<Range<Token>> ranges = new ArrayList<>(this.replicaGroups.ranges);
List<Range<Token>> ranges = new ArrayList<>(this.ranges);
ranges.sort(Range::compareTo);
return ranges;
}
@ -92,9 +108,9 @@ public class PlacementForRange
// can't use range.isWrapAround() since range.unwrap() returns a wrapping range (right token is min value)
assert range.right.compareTo(range.left) > 0 || range.right.equals(range.right.minValue());
// we're searching for an exact match to the input range here, can use standard binary search
int pos = Collections.binarySearch(replicaGroups.ranges, range, Comparator.comparing(o -> o.left));
if (pos >= 0 && pos < replicaGroups.ranges.size() && replicaGroups.ranges.get(pos).equals(range))
return replicaGroups.endpoints.get(pos);
int pos = Collections.binarySearch(ranges, range, Comparator.comparing(o -> o.left));
if (pos >= 0 && pos < ranges.size() && ranges.get(pos).equals(range))
return endpoints.get(pos);
return null;
}
@ -107,10 +123,10 @@ public class PlacementForRange
Epoch lastModified = Epoch.EMPTY;
// find a range containing the *right* token for the given range - Range is start exclusive so if we looked for the
// left one we could get the wrong range
int pos = ReplicaGroups.ordering.binarySearchAsymmetric(replicaGroups.ranges, range.right, AsymmetricOrdering.Op.CEIL);
if (pos >= 0 && pos < replicaGroups.ranges.size() && replicaGroups.ranges.get(pos).contains(range))
int pos = ordering.binarySearchAsymmetric(ranges, range.right, AsymmetricOrdering.Op.CEIL);
if (pos >= 0 && pos < ranges.size() && ranges.get(pos).contains(range))
{
VersionedEndpoints.ForRange eps = replicaGroups.endpoints.get(pos);
VersionedEndpoints.ForRange eps = endpoints.get(pos);
lastModified = eps.lastModified();
builder.addAll(eps.get(), ReplicaCollection.Builder.Conflict.ALL);
}
@ -119,10 +135,10 @@ public class PlacementForRange
public VersionedEndpoints.ForRange forRange(Token token)
{
int pos = ReplicaGroups.ordering.binarySearchAsymmetric(replicaGroups.ranges, token, AsymmetricOrdering.Op.CEIL);
if (pos >= 0 && pos < replicaGroups.endpoints.size())
return replicaGroups.endpoints.get(pos);
throw new IllegalStateException("Could not find range for token " + token + " in PlacementForRange: " + replicaGroups);
int pos = ordering.binarySearchAsymmetric(ranges, token, AsymmetricOrdering.Op.CEIL);
if (pos >= 0 && pos < endpoints.size())
return endpoints.get(pos);
throw new IllegalStateException("Could not find range for token " + token + " in PlacementForRange: " + this);
}
public VersionedEndpoints.ForToken forToken(Token token)
@ -141,8 +157,8 @@ public class PlacementForRange
public RangesByEndpoint byEndpoint()
{
RangesByEndpoint.Builder builder = new RangesByEndpoint.Builder();
for (int i = 0; i < replicaGroups.size(); i++)
replicaGroups.endpoints.get(i).byEndpoint().forEach(builder::put);
for (int i = 0; i < endpoints.size(); i++)
endpoints.get(i).byEndpoint().forEach(builder::put);
return builder.build();
}
@ -166,10 +182,10 @@ public class PlacementForRange
public PlacementForRange withCappedLastModified(Epoch lastModified)
{
SortedMap<Range<Token>, VersionedEndpoints.ForRange> copy = new TreeMap<>();
for (int i = 0; i < replicaGroups.size(); i++)
for (int i = 0; i < ranges.size(); i++)
{
Range<Token> range = replicaGroups.ranges.get(i);
VersionedEndpoints.ForRange forRange = replicaGroups.endpoints.get(i);
Range<Token> range = ranges.get(i);
VersionedEndpoints.ForRange forRange = endpoints.get(i);
if (forRange.lastModified().isAfter(lastModified))
forRange = forRange.withLastModified(lastModified);
copy.put(range, forRange);
@ -177,10 +193,38 @@ public class PlacementForRange
return new PlacementForRange(copy);
}
public int size()
{
return ranges.size();
}
public boolean isEmpty()
{
return size() == 0;
}
@VisibleForTesting
public Map<Range<Token>, VersionedEndpoints.ForRange> asMap()
{
Map<Range<Token>, VersionedEndpoints.ForRange> map = new HashMap<>();
for (int i = 0; i < size(); i++)
map.put(ranges.get(i), endpoints.get(i));
return map;
}
public void forEach(BiConsumer<Range<Token>, VersionedEndpoints.ForRange> consumer)
{
for (int i = 0; i < size(); i++)
consumer.accept(ranges.get(i), endpoints.get(i));
}
@Override
public String toString()
{
return replicaGroups.toString();
StringBuilder sb = new StringBuilder("ReplicaGroups{");
forEach((range, eps) -> sb.append(range).append('=').append(eps).append(", "));
return sb.append('}').toString();
}
@VisibleForTesting
@ -195,17 +239,16 @@ public class PlacementForRange
@VisibleForTesting
public List<String> toReplicaStringList()
{
return replicaGroups.endpoints
.stream()
.map(VersionedEndpoints.ForRange::get)
.flatMap(AbstractReplicaCollection::stream)
.map(Replica::toString)
.collect(Collectors.toList());
return endpoints.stream()
.map(VersionedEndpoints.ForRange::get)
.flatMap(AbstractReplicaCollection::stream)
.map(Replica::toString)
.collect(Collectors.toList());
}
public Builder unbuild()
{
return new Builder(replicaGroups.asMap());
return new Builder(asMap());
}
public static Builder builder()
@ -221,11 +264,11 @@ public class PlacementForRange
@VisibleForTesting
public static PlacementForRange splitRangesForPlacement(List<Token> tokens, PlacementForRange placement)
{
if (placement.replicaGroups.isEmpty())
if (placement.ranges.isEmpty())
return placement;
Builder newPlacement = PlacementForRange.builder();
List<VersionedEndpoints.ForRange> eprs = new ArrayList<>(placement.replicaGroups.endpoints);
List<VersionedEndpoints.ForRange> eprs = new ArrayList<>(placement.endpoints);
eprs.sort(Comparator.comparing(a -> a.range().left));
Token min = eprs.get(0).range().left;
Token max = eprs.get(eprs.size() - 1).range().right;
@ -393,12 +436,12 @@ public class PlacementForRange
{
public void serialize(PlacementForRange t, DataOutputPlus out, IPartitioner partitioner, Version version) throws IOException
{
out.writeInt(t.replicaGroups.size());
out.writeInt(t.ranges.size());
for (int i = 0; i < t.replicaGroups.size(); i++)
for (int i = 0; i < t.ranges.size(); i++)
{
Range<Token> range = t.replicaGroups.ranges.get(i);
VersionedEndpoints.ForRange efr = t.replicaGroups.endpoints.get(i);
Range<Token> range = t.ranges.get(i);
VersionedEndpoints.ForRange efr = t.endpoints.get(i);
if (version.isAtLeast(Version.V2))
Epoch.serializer.serialize(efr.lastModified(), out, version);
Token.metadataSerializer.serialize(range.left, out, partitioner, version);
@ -454,11 +497,11 @@ public class PlacementForRange
public long serializedSize(PlacementForRange t, IPartitioner partitioner, Version version)
{
long size = sizeof(t.replicaGroups.size());
for (int i = 0; i < t.replicaGroups.size(); i++)
long size = sizeof(t.ranges.size());
for (int i = 0; i < t.ranges.size(); i++)
{
Range<Token> range = t.replicaGroups.ranges.get(i);
VersionedEndpoints.ForRange efr = t.replicaGroups.endpoints.get(i);
Range<Token> range = t.ranges.get(i);
VersionedEndpoints.ForRange efr = t.endpoints.get(i);
if (version.isAtLeast(Version.V2))
size += Epoch.serializer.serializedSize(efr.lastModified(), version);
@ -484,93 +527,12 @@ public class PlacementForRange
if (this == o) return true;
if (!(o instanceof PlacementForRange)) return false;
PlacementForRange that = (PlacementForRange) o;
return Objects.equals(replicaGroups, that.replicaGroups);
return Objects.equals(ranges, that.ranges) && Objects.equals(endpoints, that.endpoints);
}
@Override
public int hashCode()
{
return Objects.hash(replicaGroups);
}
public static class ReplicaGroups
{
public final ImmutableList<Range<Token>> ranges;
public final ImmutableList<VersionedEndpoints.ForRange> endpoints;
private static final AsymmetricOrdering<Range<Token>, Token> ordering = new AsymmetricOrdering<>()
{
@Override
public int compare(Range<Token> left, Range<Token> right)
{
return left.compareTo(right);
}
@Override
public int compareAsymmetric(Range<Token> range, Token token)
{
if (token.isMinimum() && !range.right.isMinimum())
return -1;
if (range.left.compareTo(token) >= 0)
return 1;
if (!range.right.isMinimum() && range.right.compareTo(token) < 0)
return -1;
return 0;
}
};
public ReplicaGroups(ImmutableList<Range<Token>> ranges, ImmutableList<VersionedEndpoints.ForRange> endpoints)
{
this.ranges = ranges;
this.endpoints = endpoints;
}
public int size()
{
return ranges.size();
}
public boolean isEmpty()
{
return size() == 0;
}
@VisibleForTesting
public Map<Range<Token>, VersionedEndpoints.ForRange> asMap()
{
Map<Range<Token>, VersionedEndpoints.ForRange> map = new HashMap<>();
for (int i = 0; i < size(); i++)
map.put(ranges.get(i), endpoints.get(i));
return map;
}
public void forEach(BiConsumer<Range<Token>, VersionedEndpoints.ForRange> consumer)
{
for (int i = 0; i < size(); i++)
consumer.accept(ranges.get(i), endpoints.get(i));
}
@Override
public String toString()
{
StringBuilder sb = new StringBuilder("ReplicaGroups{");
forEach((range, eps) -> sb.append(range).append('=').append(eps).append(", "));
return sb.append('}').toString();
}
@Override
public boolean equals(Object o)
{
if (this == o) return true;
if (!(o instanceof ReplicaGroups)) return false;
ReplicaGroups entries = (ReplicaGroups) o;
return Objects.equals(ranges, entries.ranges) && Objects.equals(endpoints, entries.endpoints);
}
@Override
public int hashCode()
{
return Objects.hash(ranges, endpoints);
}
return Objects.hash(ranges, endpoints);
}
}

View File

@ -558,7 +558,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.replicaGroups().endpoints)
for (VersionedEndpoints.ForRange placements : sut.service.metadata().placements.get(rf.asKeyspaceParams().replication).writes.endpoints)
{
List<NodeId> replicas = new ArrayList<>(metadata.directory.toNodeIds(placements.get().endpoints()));
List<NodeId> bounceCandidates = new ArrayList<>();
@ -702,7 +702,7 @@ public class MetadataChangeSimulationTest extends CMSTestBase
public static void match(PlacementForRange actual, Map<TokenPlacementModel.Range, List<TokenPlacementModel.Replica>> predicted) throws Throwable
{
Map<Range<Token>, VersionedEndpoints.ForRange> actualGroups = actual.replicaGroups().asMap();
Map<Range<Token>, VersionedEndpoints.ForRange> actualGroups = actual.asMap();
assert predicted.size() == actualGroups.size() :
String.format("\nPredicted:\n%s(%d)" +
"\nActual:\n%s(%d)", toString(predicted), predicted.size(), toString(actualGroups), actualGroups.size());

View File

@ -133,14 +133,14 @@ public class OperationalEquivalenceTest extends CMSTestBase
{
l.forEach((params, lPlacement) -> {
DataPlacement rPlacement = r.get(params);
lPlacement.reads.replicaGroups().forEach((range, lReplicas) -> {
lPlacement.reads.forEach((range, lReplicas) -> {
EndpointsForRange rReplicas = rPlacement.reads.forRange(range).get();
Assert.assertEquals(toReplicas(lReplicas.get()),
toReplicas(rReplicas));
});
lPlacement.writes.replicaGroups().forEach((range, lReplicas) -> {
lPlacement.writes.forEach((range, lReplicas) -> {
EndpointsForRange rReplicas = rPlacement.writes.forRange(range).get();
Assert.assertEquals(toReplicas(lReplicas.get()),

View File

@ -156,8 +156,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.replicaGroups().asMap();
m1.placements.get(simulatedState.rf.asKeyspaceParams().replication).reads.replicaGroups().forEach((k, beforePlacements) -> {
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) -> {
if (after.containsKey(k))
{
VersionedEndpoints.ForRange afterPlacements = after.get(k);

View File

@ -54,9 +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.replicaGroups().endpoints)
for (VersionedEndpoints.ForRange fr : metadata.placements.get(ReplicationParams.simple(i)).writes.endpoints)
{
if (smallestSeen == null || fr.lastModified().isBefore(smallestSeen))
smallestSeen = fr.lastModified();

View File

@ -266,7 +266,7 @@ public class UniformRangePlacementTest
DataPlacement initialPlacement = builder.build();
List<Token> tokens = ImmutableList.of(token(Long.MIN_VALUE), token(0));
DataPlacement newPlacement = initialPlacement.splitRangesForPlacement(tokens);
assertEquals(2, newPlacement.writes.replicaGroups().size());
assertEquals(2, newPlacement.writes.size());
}
private PlacementForRange initialPlacement()
@ -276,7 +276,7 @@ public class UniformRangePlacementTest
rg(200, 300, 1, 2, 3),
rg(300, 400, 1, 2, 3) };
PlacementForRange placement = PlacementForRange.builder()
.withReplicaGroups(Arrays.asList(initialGroups).stream().map(this::v).collect(Collectors.toList()))
.withReplicaGroups(Arrays.stream(initialGroups).map(this::v).collect(Collectors.toList()))
.build();
assertPlacement(placement, initialGroups);
return placement;
@ -284,7 +284,7 @@ public class UniformRangePlacementTest
private void assertPlacement(PlacementForRange placement, EndpointsForRange...expected)
{
Collection<EndpointsForRange> replicaGroups = placement.replicaGroups().endpoints.stream().map(v -> v.get()).collect(Collectors.toList());
Collection<EndpointsForRange> replicaGroups = placement.endpoints.stream().map(VersionedEndpoints.ForRange::get).collect(Collectors.toList());
assertEquals(replicaGroups.size(), expected.length);
int i = 0;
boolean allMatch = true;

View File

@ -323,8 +323,8 @@ public class InProgressSequenceCancellationTest
DataPlacement otherPlacement = second.get(params);
PlacementForRange r1 = placement.reads;
PlacementForRange r2 = otherPlacement.reads;
assertEquals(r1.replicaGroups().ranges, r2.replicaGroups().ranges);
r1.replicaGroups().forEach((range, e1) -> {
assertEquals(r1.ranges, r2.ranges);
r1.forEach((range, e1) -> {
EndpointsForRange e2 = r2.forRange(range).get();
assertEquals(e1.size(),e2.size());
assertTrue(e1.get().stream().allMatch(e2::contains));
@ -332,8 +332,8 @@ public class InProgressSequenceCancellationTest
PlacementForRange w1 = placement.reads;
PlacementForRange w2 = otherPlacement.reads;
assertEquals(w1.replicaGroups().ranges, w2.replicaGroups().ranges);
w1.replicaGroups().forEach((range, e1) -> {
assertEquals(w1.ranges, w2.ranges);
w1.forEach((range, e1) -> {
EndpointsForRange e2 = w2.forRange(range).get();
assertEquals(e1.size(),e2.size());
assertTrue(e1.get().stream().allMatch(e2::contains));

View File

@ -67,7 +67,7 @@ public class SequencesUtils
{
LockedRanges.AffectedRangesBuilder affected = LockedRanges.AffectedRanges.builder();
placements.asMap().forEach((params, placement) -> {
placement.reads.replicaGroups().ranges.forEach((range) -> {
placement.reads.ranges.forEach((range) -> {
if (random.nextDouble() >= 0.6)
affected.add(params, range);
});