mirror of https://github.com/apache/cassandra
Properly set lastModifiedEpoch on multistep operations
Patch by Marcus Eriksson; reviewed by Sam Tunnicliffe for CASSANDRA-19538
This commit is contained in:
parent
d9192745bc
commit
80971709b9
|
|
@ -155,7 +155,7 @@ public abstract class MultiStepOperation<CONTEXT>
|
||||||
public static Transformation.Result applyMultipleTransformations(ClusterMetadata metadata, Transformation.Kind next, List<Transformation> transformations)
|
public static Transformation.Result applyMultipleTransformations(ClusterMetadata metadata, Transformation.Kind next, List<Transformation> transformations)
|
||||||
{
|
{
|
||||||
ImmutableSet.Builder<MetadataKey> modifiedKeys = ImmutableSet.builder();
|
ImmutableSet.Builder<MetadataKey> modifiedKeys = ImmutableSet.builder();
|
||||||
Epoch lastModifiedEpoch = metadata.epoch;
|
Epoch lastModifiedEpoch = metadata.epoch.nextEpoch();
|
||||||
boolean foundStart = false;
|
boolean foundStart = false;
|
||||||
for (Transformation nextTransformation : transformations)
|
for (Transformation nextTransformation : transformations)
|
||||||
{
|
{
|
||||||
|
|
@ -169,7 +169,7 @@ public abstract class MultiStepOperation<CONTEXT>
|
||||||
modifiedKeys.addAll(result.success().affectedMetadata);
|
modifiedKeys.addAll(result.success().affectedMetadata);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return new Transformation.Success(metadata.forceEpoch(lastModifiedEpoch.nextEpoch()), LockedRanges.AffectedRanges.EMPTY, modifiedKeys.build());
|
return new Transformation.Success(metadata, LockedRanges.AffectedRanges.EMPTY, modifiedKeys.build());
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
|
||||||
|
|
@ -155,6 +155,12 @@ public class DataPlacement
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public DataPlacement withCappedLastModified(Epoch lastModified)
|
||||||
|
{
|
||||||
|
return new DataPlacement(reads.withCappedLastModified(lastModified),
|
||||||
|
writes.withCappedLastModified(lastModified));
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public String toString()
|
public String toString()
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -25,6 +25,7 @@ import java.util.Map;
|
||||||
import java.util.Objects;
|
import java.util.Objects;
|
||||||
import java.util.function.BiConsumer;
|
import java.util.function.BiConsumer;
|
||||||
|
|
||||||
|
import com.google.common.collect.ImmutableMap;
|
||||||
import com.google.common.collect.Maps;
|
import com.google.common.collect.Maps;
|
||||||
import com.google.common.collect.Sets;
|
import com.google.common.collect.Sets;
|
||||||
|
|
||||||
|
|
@ -122,7 +123,7 @@ public class DataPlacements extends ReplicationMap<DataPlacement> implements Met
|
||||||
@Override
|
@Override
|
||||||
public DataPlacements withLastModified(Epoch epoch)
|
public DataPlacements withLastModified(Epoch epoch)
|
||||||
{
|
{
|
||||||
return new DataPlacements(epoch, asMap());
|
return new DataPlacements(epoch, capLastModified(epoch, map));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -255,6 +256,13 @@ public class DataPlacements extends ReplicationMap<DataPlacement> implements Met
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public static ImmutableMap<ReplicationParams, DataPlacement> capLastModified(Epoch lastModified, Map<ReplicationParams, DataPlacement> placements)
|
||||||
|
{
|
||||||
|
ImmutableMap.Builder<ReplicationParams, DataPlacement> builder = ImmutableMap.builder();
|
||||||
|
placements.forEach((params, placement) -> builder.put(params, placement.withCappedLastModified(lastModified)));
|
||||||
|
return builder.build();
|
||||||
|
}
|
||||||
|
|
||||||
public void dumpDiff(DataPlacements other)
|
public void dumpDiff(DataPlacements other)
|
||||||
{
|
{
|
||||||
if (!map.equals(other.map))
|
if (!map.equals(other.map))
|
||||||
|
|
|
||||||
|
|
@ -152,6 +152,20 @@ public class PlacementForRange
|
||||||
return builder.build();
|
return builder.build();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public PlacementForRange withCappedLastModified(Epoch lastModified)
|
||||||
|
{
|
||||||
|
SortedMap<Range<Token>, VersionedEndpoints.ForRange> copy = new TreeMap<>();
|
||||||
|
for (Map.Entry<Range<Token>, VersionedEndpoints.ForRange> entry : replicaGroups.entrySet())
|
||||||
|
{
|
||||||
|
Range<Token> range = entry.getKey();
|
||||||
|
VersionedEndpoints.ForRange forRange = entry.getValue();
|
||||||
|
if (forRange.lastModified().isAfter(lastModified))
|
||||||
|
forRange = forRange.withLastModified(lastModified);
|
||||||
|
copy.put(range, forRange);
|
||||||
|
}
|
||||||
|
return new PlacementForRange(copy);
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public String toString()
|
public String toString()
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -63,7 +63,7 @@ public interface VersionedEndpoints<E extends Endpoints<E>> extends MetadataValu
|
||||||
|
|
||||||
public ForRange withLastModified(Epoch epoch)
|
public ForRange withLastModified(Epoch epoch)
|
||||||
{
|
{
|
||||||
return new ForRange(lastModified, endpointsForRange);
|
return new ForRange(epoch, endpointsForRange);
|
||||||
}
|
}
|
||||||
|
|
||||||
public Epoch lastModified()
|
public Epoch lastModified()
|
||||||
|
|
@ -151,7 +151,7 @@ public interface VersionedEndpoints<E extends Endpoints<E>> extends MetadataValu
|
||||||
|
|
||||||
public ForToken withLastModified(Epoch epoch)
|
public ForToken withLastModified(Epoch epoch)
|
||||||
{
|
{
|
||||||
return new ForToken(lastModified, endpointsForToken);
|
return new ForToken(epoch, endpointsForToken);
|
||||||
}
|
}
|
||||||
|
|
||||||
public ForToken map(Function<EndpointsForToken, EndpointsForToken> fn)
|
public ForToken map(Function<EndpointsForToken, EndpointsForToken> fn)
|
||||||
|
|
@ -159,13 +159,11 @@ public interface VersionedEndpoints<E extends Endpoints<E>> extends MetadataValu
|
||||||
return new ForToken(lastModified, fn.apply(endpointsForToken));
|
return new ForToken(lastModified, fn.apply(endpointsForToken));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
public ForToken without(Set<InetAddressAndPort> remove)
|
public ForToken without(Set<InetAddressAndPort> remove)
|
||||||
{
|
{
|
||||||
return map(e -> e.without(remove));
|
return map(e -> e.without(remove));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
public Epoch lastModified()
|
public Epoch lastModified()
|
||||||
{
|
{
|
||||||
return lastModified;
|
return lastModified;
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue