CASSANDRA-19191

This commit is contained in:
Marcus Eriksson 2024-03-20 15:53:50 +01:00
parent c5c4cd4e57
commit 6af507a11c
9 changed files with 155 additions and 77 deletions

View File

@ -21,7 +21,6 @@ package org.apache.cassandra.tcm.ownership;
import java.io.IOException;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
@ -78,12 +77,12 @@ public class DataPlacement
public DataPlacement combineReplicaGroups(DataPlacement other)
{
return new DataPlacement(PlacementForRange.builder()
.withReplicaGroups(reads.replicaGroups().values())
.withReplicaGroups(other.reads.replicaGroups.values())
.withReplicaGroups(reads.replicaGroups().endpoints)
.withReplicaGroups(other.reads.replicaGroups().endpoints)
.build(),
PlacementForRange.builder()
.withReplicaGroups(writes.replicaGroups().values())
.withReplicaGroups(other.writes.replicaGroups.values())
.withReplicaGroups(writes.replicaGroups().endpoints)
.withReplicaGroups(other.writes.replicaGroups().endpoints)
.build());
}
@ -165,22 +164,11 @@ public class DataPlacement
public String toString()
{
return "DataPlacement{" +
"reads=" + toString(reads.replicaGroups) +
", writes=" + toString(writes.replicaGroups) +
"reads=" + reads +
", writes=" + writes +
'}';
}
public static String toString(Map<?, ?> predicted)
{
StringBuilder sb = new StringBuilder();
for (Map.Entry<?, ?> e : predicted.entrySet())
{
sb.append(e.getKey()).append("=").append(e.getValue()).append(",\n");
}
return sb.toString();
}
@Override
public boolean equals(Object o)
{

View File

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

View File

@ -20,9 +20,12 @@ package org.apache.cassandra.tcm.ownership;
import java.io.IOException;
import java.util.*;
import java.util.function.BiConsumer;
import java.util.stream.Collectors;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableSortedMap;
import com.google.common.collect.Maps;
import org.apache.cassandra.dht.IPartitioner;
@ -41,6 +44,7 @@ import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.tcm.Epoch;
import org.apache.cassandra.tcm.serialization.PartitionerAwareMetadataSerializer;
import org.apache.cassandra.tcm.serialization.Version;
import org.apache.cassandra.utils.AsymmetricOrdering;
import static org.apache.cassandra.db.TypeSizes.sizeof;
@ -50,23 +54,34 @@ public class PlacementForRange
public static final PlacementForRange EMPTY = PlacementForRange.builder().build();
final SortedMap<Range<Token>, VersionedEndpoints.ForRange> replicaGroups;
private final ReplicaGroups replicaGroups;
public PlacementForRange(Map<Range<Token>, VersionedEndpoints.ForRange> replicaGroups)
{
this.replicaGroups = new TreeMap<>(replicaGroups);
ImmutableList.Builder<Range<Token>> rangesBuilder = ImmutableList.builderWithExpectedSize(replicaGroups.size());
ImmutableList.Builder<VersionedEndpoints.ForRange> endpointsBuilder = ImmutableList.builderWithExpectedSize(replicaGroups.size());
Range<Token> prev = null;
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> entry : ImmutableSortedMap.copyOf(replicaGroups, Comparator.comparing(o -> o.left)).entrySet())
{
if (prev != null && prev.right.compareTo(entry.getKey().left) > 0 )
throw new IllegalArgumentException("Got overlapping ranges in replica groups: " + replicaGroups);
prev = entry.getKey();
rangesBuilder.add(entry.getKey());
endpointsBuilder.add(entry.getValue());
}
this.replicaGroups = new ReplicaGroups(rangesBuilder.build(), endpointsBuilder.build());
}
@VisibleForTesting
public Map<Range<Token>, VersionedEndpoints.ForRange> replicaGroups()
public ReplicaGroups replicaGroups()
{
return Collections.unmodifiableMap(replicaGroups);
return replicaGroups;
}
@VisibleForTesting
public List<Range<Token>> ranges()
{
List<Range<Token>> ranges = new ArrayList<>(replicaGroups.keySet());
List<Range<Token>> ranges = new ArrayList<>(this.replicaGroups.ranges);
ranges.sort(Range::compareTo);
return ranges;
}
@ -76,7 +91,11 @@ 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());
return replicaGroups.get(range);
// 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);
return null;
}
/**
@ -86,37 +105,29 @@ public class PlacementForRange
{
EndpointsForRange.Builder builder = new EndpointsForRange.Builder(range);
Epoch lastModified = Epoch.EMPTY;
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> entry : replicaGroups.entrySet())
// 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))
{
if (entry.getKey().contains(range))
{
lastModified = Epoch.max(lastModified, entry.getValue().lastModified());
builder.addAll(entry.getValue().get(), ReplicaCollection.Builder.Conflict.ALL);
}
VersionedEndpoints.ForRange eps = replicaGroups.endpoints.get(pos);
lastModified = eps.lastModified();
builder.addAll(eps.get(), ReplicaCollection.Builder.Conflict.ALL);
}
return VersionedEndpoints.forRange(lastModified, builder.build());
}
public VersionedEndpoints.ForRange forRange(Token token)
{
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> entry : replicaGroups.entrySet())
{
if (entry.getKey().contains(token))
return entry.getValue();
}
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);
}
public VersionedEndpoints.ForToken forToken(Token token)
{
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> entry : replicaGroups.entrySet())
{
if (entry.getKey().contains(token))
return entry.getValue().forToken(token);
}
throw new IllegalStateException("Could not find range for token " + token + " in PlacementForRange: " + replicaGroups);
return forRange(token).forToken(token);
}
public Delta difference(PlacementForRange next)
@ -130,8 +141,8 @@ public class PlacementForRange
public RangesByEndpoint byEndpoint()
{
RangesByEndpoint.Builder builder = new RangesByEndpoint.Builder();
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> oldPlacement : this.replicaGroups.entrySet())
oldPlacement.getValue().byEndpoint().forEach(builder::put);
for (int i = 0; i < replicaGroups.size(); i++)
replicaGroups.endpoints.get(i).byEndpoint().forEach(builder::put);
return builder.build();
}
@ -155,10 +166,10 @@ public class PlacementForRange
public PlacementForRange withCappedLastModified(Epoch lastModified)
{
SortedMap<Range<Token>, VersionedEndpoints.ForRange> copy = new TreeMap<>();
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> entry : replicaGroups.entrySet())
for (int i = 0; i < replicaGroups.size(); i++)
{
Range<Token> range = entry.getKey();
VersionedEndpoints.ForRange forRange = entry.getValue();
Range<Token> range = replicaGroups.ranges.get(i);
VersionedEndpoints.ForRange forRange = replicaGroups.endpoints.get(i);
if (forRange.lastModified().isAfter(lastModified))
forRange = forRange.withLastModified(lastModified);
copy.put(range, forRange);
@ -184,7 +195,7 @@ public class PlacementForRange
@VisibleForTesting
public List<String> toReplicaStringList()
{
return replicaGroups.values()
return replicaGroups.endpoints
.stream()
.map(VersionedEndpoints.ForRange::get)
.flatMap(AbstractReplicaCollection::stream)
@ -194,7 +205,7 @@ public class PlacementForRange
public Builder unbuild()
{
return new Builder(new HashMap<>(replicaGroups));
return new Builder(replicaGroups.asMap());
}
public static Builder builder()
@ -214,7 +225,7 @@ public class PlacementForRange
return placement;
Builder newPlacement = PlacementForRange.builder();
List<VersionedEndpoints.ForRange> eprs = new ArrayList<>(placement.replicaGroups.values());
List<VersionedEndpoints.ForRange> eprs = new ArrayList<>(placement.replicaGroups.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;
@ -384,18 +395,18 @@ public class PlacementForRange
{
out.writeInt(t.replicaGroups.size());
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> entry : t.replicaGroups.entrySet())
for (int i = 0; i < t.replicaGroups.size(); i++)
{
Range<Token> range = entry.getKey();
VersionedEndpoints.ForRange efr = entry.getValue();
Range<Token> range = t.replicaGroups.ranges.get(i);
VersionedEndpoints.ForRange efr = t.replicaGroups.endpoints.get(i);
if (version.isAtLeast(Version.V2))
Epoch.serializer.serialize(efr.lastModified(), out, version);
Token.metadataSerializer.serialize(range.left, out, partitioner, version);
Token.metadataSerializer.serialize(range.right, out, partitioner, version);
out.writeInt(efr.size());
for (int i = 0; i < efr.size(); i++)
for (int efrIdx = 0; efrIdx < efr.size(); efrIdx++)
{
Replica r = efr.get().get(i);
Replica r = efr.get().get(efrIdx);
Token.metadataSerializer.serialize(r.range().left, out, partitioner, version);
Token.metadataSerializer.serialize(r.range().right, out, partitioner, version);
InetAddressAndPort.MetadataSerializer.serializer.serialize(r.endpoint(), out, version);
@ -444,19 +455,19 @@ public class PlacementForRange
public long serializedSize(PlacementForRange t, IPartitioner partitioner, Version version)
{
long size = sizeof(t.replicaGroups.size());
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> entry : t.replicaGroups.entrySet())
for (int i = 0; i < t.replicaGroups.size(); i++)
{
Range<Token> range = entry.getKey();
VersionedEndpoints.ForRange efr = entry.getValue();
Range<Token> range = t.replicaGroups.ranges.get(i);
VersionedEndpoints.ForRange efr = t.replicaGroups.endpoints.get(i);
if (version.isAtLeast(Version.V2))
size += Epoch.serializer.serializedSize(efr.lastModified(), version);
size += Token.metadataSerializer.serializedSize(range.left, partitioner, version);
size += Token.metadataSerializer.serializedSize(range.right, partitioner, version);
size += sizeof(efr.size());
for (int i = 0; i < efr.size(); i++)
for (int efrIdx = 0; efrIdx < efr.size(); efrIdx++)
{
Replica r = efr.get().get(i);
Replica r = efr.get().get(efrIdx);
size += Token.metadataSerializer.serializedSize(r.range().left, partitioner, version);
size += Token.metadataSerializer.serializedSize(r.range().right, partitioner, version);
size += InetAddressAndPort.MetadataSerializer.serializer.serializedSize(r.endpoint(), version);
@ -481,4 +492,85 @@ public class PlacementForRange
{
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);
}
}
}

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().values())
for (VersionedEndpoints.ForRange placements : sut.service.metadata().placements.get(rf.asKeyspaceParams().replication).writes.replicaGroups().endpoints)
{
List<NodeId> replicas = new ArrayList<>(metadata.directory.toNodeIds(placements.get().endpoints()));
List<NodeId> bounceCandidates = new ArrayList<>();
@ -702,10 +702,10 @@ 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();
Map<Range<Token>, VersionedEndpoints.ForRange> actualGroups = actual.replicaGroups().asMap();
assert predicted.size() == actualGroups.size() :
String.format("\nPredicted:\n%s(%d)" +
"\nActual:\n%s(%d)", toString(predicted), predicted.size(), toString(actual.replicaGroups()), actualGroups.size());
"\nActual:\n%s(%d)", toString(predicted), predicted.size(), toString(actualGroups), actualGroups.size());
for (Map.Entry<TokenPlacementModel.Range, List<TokenPlacementModel.Replica>> entry : predicted.entrySet())
{

View File

@ -156,7 +156,7 @@ 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();
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) -> {
if (after.containsKey(k))
{

View File

@ -56,7 +56,7 @@ public class RangeVersioningTest extends FuzzTestBase
Epoch smallestSeen = null;
for (VersionedEndpoints.ForRange fr : metadata.placements
.get(ReplicationParams.simple(i))
.writes.replicaGroups().values())
.writes.replicaGroups().endpoints)
{
if (smallestSeen == null || fr.lastModified().isBefore(smallestSeen))
smallestSeen = fr.lastModified();

View File

@ -235,9 +235,7 @@ public class UniformRangePlacementTest
DataPlacement initialPlacement = builder.build();
DataPlacement split = initialPlacement.splitRangesForPlacement(tokens);
assertPlacement(split.writes, rg(-3, -9223372036854775808L, 1),
rg(-9223372036854775808L,-4611686018427387905L, 1),
rg(-4611686018427387905L, -3, 1));
assertPlacement(split.writes, rg(-9223372036854775808L,-4611686018427387905L, 1), rg(-4611686018427387905L, -3, 1), rg(-3, -9223372036854775808L, 1));
}
@Test
@ -254,7 +252,7 @@ public class UniformRangePlacementTest
DataPlacement initialPlacement = builder.build();
DataPlacement split = initialPlacement.splitRangesForPlacement(tokens);
assertPlacement(split.writes, rg(3074457345618258602L,-9223372036854775808L, 1), rg(-9223372036854775808L, 3074457345618258602L, 1));
assertPlacement(split.writes, rg(-9223372036854775808L, 3074457345618258602L, 1), rg(3074457345618258602L,-9223372036854775808L, 1));
}
@Test
@ -268,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.values().size());
assertEquals(2, newPlacement.writes.replicaGroups().size());
}
private PlacementForRange initialPlacement()
@ -286,7 +284,7 @@ public class UniformRangePlacementTest
private void assertPlacement(PlacementForRange placement, EndpointsForRange...expected)
{
Collection<EndpointsForRange> replicaGroups = placement.replicaGroups.values().stream().map(v -> v.get()).collect(Collectors.toList());
Collection<EndpointsForRange> replicaGroups = placement.replicaGroups().endpoints.stream().map(v -> v.get()).collect(Collectors.toList());
assertEquals(replicaGroups.size(), expected.length);
int i = 0;
boolean allMatch = true;

View File

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

View File

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