mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-5.0' into trunk
* cassandra-5.0: CASSANDRA-19633 Replaced node is stuck in a loop calculating ranges
This commit is contained in:
commit
f7462446a4
|
|
@ -494,6 +494,13 @@ public enum CassandraRelevantProperties
|
|||
SIZE_RECORDER_INTERVAL("cassandra.size_recorder_interval", "300"),
|
||||
SKIP_AUTH_SETUP("cassandra.skip_auth_setup", "false"),
|
||||
SKIP_GC_INSPECTOR("cassandra.skip_gc_inspector", "false"),
|
||||
|
||||
/**
|
||||
* Do not try to calculate optimal streaming candidates. This can take a lot of time in some configs specially
|
||||
* with vnodes.
|
||||
*/
|
||||
SKIP_OPTIMAL_STREAMING_CANDIDATES_CALCULATION("cassandra.skip_optimal_streaming_candidates_calculation", "false"),
|
||||
|
||||
SKIP_PAXOS_REPAIR_ON_TOPOLOGY_CHANGE("cassandra.skip_paxos_repair_on_topology_change"),
|
||||
/** If necessary for operational purposes, permit certain keyspaces to be ignored for paxos topology repairs. */
|
||||
SKIP_PAXOS_REPAIR_ON_TOPOLOGY_CHANGE_KEYSPACES("cassandra.skip_paxos_repair_on_topology_change_keyspaces"),
|
||||
|
|
|
|||
|
|
@ -39,6 +39,7 @@ import org.apache.commons.lang3.StringUtils;
|
|||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import org.apache.cassandra.config.CassandraRelevantProperties;
|
||||
import org.apache.cassandra.db.Keyspace;
|
||||
import org.apache.cassandra.db.SystemKeyspace;
|
||||
import org.apache.cassandra.gms.FailureDetector;
|
||||
|
|
@ -383,8 +384,12 @@ public class RangeStreamer
|
|||
|
||||
Multimap<InetAddressAndPort, FetchReplica> workMap;
|
||||
//Only use the optimized strategy if we don't care about strict sources, have a replication factor > 1, and no
|
||||
//transient replicas.
|
||||
if (useStrictSource || strat == null || strat.getReplicationFactor().allReplicas == 1 || strat.getReplicationFactor().hasTransientReplicas())
|
||||
//transient replicas or it is intentionally skipped.
|
||||
if (CassandraRelevantProperties.SKIP_OPTIMAL_STREAMING_CANDIDATES_CALCULATION.getBoolean() ||
|
||||
useStrictSource ||
|
||||
strat == null ||
|
||||
strat.getReplicationFactor().allReplicas == 1 ||
|
||||
strat.getReplicationFactor().hasTransientReplicas())
|
||||
{
|
||||
workMap = convertPreferredEndpointsToWorkMap(fetchMap);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -23,6 +23,8 @@ import java.util.List;
|
|||
import java.util.Random;
|
||||
import java.util.Set;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
|
||||
import com.google.common.base.Predicate;
|
||||
import com.google.common.base.Predicates;
|
||||
|
|
@ -31,9 +33,11 @@ import com.google.common.collect.Multimap;
|
|||
import org.junit.AfterClass;
|
||||
import org.junit.BeforeClass;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
|
||||
import org.apache.cassandra.SchemaLoader;
|
||||
import org.apache.cassandra.ServerTestUtils;
|
||||
import org.apache.cassandra.config.CassandraRelevantProperties;
|
||||
import org.apache.cassandra.config.DatabaseDescriptor;
|
||||
import org.apache.cassandra.db.Keyspace;
|
||||
import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper;
|
||||
|
|
@ -51,16 +55,35 @@ import org.apache.cassandra.tcm.membership.NodeId;
|
|||
import org.apache.cassandra.tcm.ownership.MovementMap;
|
||||
import org.apache.cassandra.tcm.sequences.BootstrapAndJoin;
|
||||
import org.apache.cassandra.utils.Pair;
|
||||
import org.jboss.byteman.contrib.bmunit.BMRule;
|
||||
import org.jboss.byteman.contrib.bmunit.BMRules;
|
||||
import org.jboss.byteman.contrib.bmunit.BMUnitRunner;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
|
||||
@RunWith(BMUnitRunner.class)
|
||||
public class BootStrapperTest
|
||||
{
|
||||
static IPartitioner oldPartitioner;
|
||||
static Predicate<Replica> originalAlivePredicate = RangeStreamer.ALIVE_PREDICATE;
|
||||
public static AtomicBoolean nonOptimizationHit = new AtomicBoolean(false);
|
||||
public static AtomicBoolean optimizationHit = new AtomicBoolean(false);
|
||||
private static final IFailureDetector mockFailureDetector = new IFailureDetector()
|
||||
{
|
||||
public boolean isAlive(InetAddressAndPort ep)
|
||||
{
|
||||
return true;
|
||||
}
|
||||
|
||||
public void interpret(InetAddressAndPort ep) { throw new UnsupportedOperationException(); }
|
||||
public void report(InetAddressAndPort ep) { throw new UnsupportedOperationException(); }
|
||||
public void registerFailureDetectionEventListener(IFailureDetectionEventListener listener) { throw new UnsupportedOperationException(); }
|
||||
public void unregisterFailureDetectionEventListener(IFailureDetectionEventListener listener) { throw new UnsupportedOperationException(); }
|
||||
public void remove(InetAddressAndPort ep) { throw new UnsupportedOperationException(); }
|
||||
public void forceConviction(InetAddressAndPort ep) { throw new UnsupportedOperationException(); }
|
||||
};
|
||||
|
||||
@BeforeClass
|
||||
public static void setup() throws ConfigurationException
|
||||
|
|
@ -96,46 +119,49 @@ public class BootStrapperTest
|
|||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@BMRules(rules = { @BMRule(name = "Make sure the non-optimized path is picked up for some operations",
|
||||
targetClass = "org.apache.cassandra.dht.RangeStreamer",
|
||||
targetMethod = "convertPreferredEndpointsToWorkMap(EndpointsByReplica)",
|
||||
action = "org.apache.cassandra.dht.BootStrapperTest.nonOptimizationHit.set(true)"),
|
||||
@BMRule(name = "Make sure the optimized path is picked up for some operations",
|
||||
targetClass = "org.apache.cassandra.dht.RangeStreamer",
|
||||
targetMethod = "getOptimizedWorkMap(EndpointsByReplica,Collection,String)",
|
||||
action = "org.apache.cassandra.dht.BootStrapperTest.optimizationHit.set(true)") })
|
||||
public void testStreamingCandidatesOptmizationSkip() throws UnknownHostException
|
||||
{
|
||||
testSkipStreamingCandidatesOptmizationFeatureFlag(true, true, false, getRangeStreamer());
|
||||
testSkipStreamingCandidatesOptmizationFeatureFlag(false, true, true, getRangeStreamer());
|
||||
}
|
||||
|
||||
private void testSkipStreamingCandidatesOptmizationFeatureFlag(boolean disableOptimization, boolean nonOptimizedPathHit, boolean optimizedPathHit, RangeStreamer s) throws UnknownHostException
|
||||
{
|
||||
try
|
||||
{
|
||||
nonOptimizationHit.set(false);
|
||||
optimizationHit.set(false);
|
||||
CassandraRelevantProperties.SKIP_OPTIMAL_STREAMING_CANDIDATES_CALCULATION.setBoolean(disableOptimization);
|
||||
|
||||
for (String keyspaceName : Schema.instance.getUserKeyspaces().names())
|
||||
s.addKeyspaceToFetch(keyspaceName);
|
||||
|
||||
assertEquals(nonOptimizedPathHit, nonOptimizationHit.get());
|
||||
if (disableOptimization) // The optimized path may or not be hit depending on the code.
|
||||
assertEquals(optimizedPathHit, optimizationHit.get());
|
||||
}
|
||||
finally
|
||||
{
|
||||
CassandraRelevantProperties.SKIP_OPTIMAL_STREAMING_CANDIDATES_CALCULATION.reset();
|
||||
}
|
||||
}
|
||||
|
||||
private RangeStreamer testSourceTargetComputation(String keyspaceName, int numOldNodes, int replicationFactor) throws UnknownHostException
|
||||
{
|
||||
ServerTestUtils.resetCMS();
|
||||
generateFakeEndpoints(numOldNodes);
|
||||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
|
||||
assertEquals(numOldNodes, metadata.tokenMap.tokens().size());
|
||||
IFailureDetector mockFailureDetector = new IFailureDetector()
|
||||
{
|
||||
public boolean isAlive(InetAddressAndPort ep)
|
||||
{
|
||||
return true;
|
||||
}
|
||||
|
||||
public void interpret(InetAddressAndPort ep) { throw new UnsupportedOperationException(); }
|
||||
public void report(InetAddressAndPort ep) { throw new UnsupportedOperationException(); }
|
||||
public void registerFailureDetectionEventListener(IFailureDetectionEventListener listener) { throw new UnsupportedOperationException(); }
|
||||
public void unregisterFailureDetectionEventListener(IFailureDetectionEventListener listener) { throw new UnsupportedOperationException(); }
|
||||
public void remove(InetAddressAndPort ep) { throw new UnsupportedOperationException(); }
|
||||
public void forceConviction(InetAddressAndPort ep) { throw new UnsupportedOperationException(); }
|
||||
};
|
||||
|
||||
Token myToken = metadata.partitioner.getRandomToken();
|
||||
InetAddressAndPort myEndpoint = InetAddressAndPort.getByName("127.0.0.1");
|
||||
NodeId newNode = ClusterMetadataTestHelper.register(myEndpoint);
|
||||
ClusterMetadataTestHelper.JoinProcess join = ClusterMetadataTestHelper.lazyJoin(myEndpoint, myToken);
|
||||
join.prepareJoin();
|
||||
metadata = ClusterMetadata.current();
|
||||
BootstrapAndJoin joiningPlan = (BootstrapAndJoin) metadata.inProgressSequences.get(newNode);
|
||||
Pair<MovementMap, MovementMap> movements = joiningPlan.getMovementMaps(metadata);
|
||||
RangeStreamer s = new RangeStreamer(metadata,
|
||||
StreamOperation.BOOTSTRAP,
|
||||
true,
|
||||
DatabaseDescriptor.getNodeProximity(),
|
||||
new StreamStateStore(),
|
||||
mockFailureDetector,
|
||||
false,
|
||||
1,
|
||||
movements.left,
|
||||
movements.right);
|
||||
RangeStreamer s = getRangeStreamer();
|
||||
|
||||
assertNotNull(Keyspace.open(keyspaceName));
|
||||
s.addKeyspaceToFetch(keyspaceName);
|
||||
|
|
@ -159,10 +185,40 @@ public class BootStrapperTest
|
|||
// there isn't any point in testing the size of these collections for any specific size. When a random partitioner
|
||||
// is used, they will vary.
|
||||
assert toFetch.values().size() > 0;
|
||||
assert toFetch.keys().stream().noneMatch(myEndpoint::equals);
|
||||
|
||||
assert toFetch.keys().stream().noneMatch(InetAddressAndPort.getByName("127.0.0.1")::equals);
|
||||
return s;
|
||||
}
|
||||
|
||||
private RangeStreamer getRangeStreamer() throws UnknownHostException
|
||||
{
|
||||
ClusterMetadata metadata = ClusterMetadata.current();
|
||||
Pair<MovementMap, MovementMap> movements = Pair.create(MovementMap.empty(), MovementMap.empty());
|
||||
|
||||
if (metadata.myNodeId() == null)
|
||||
{
|
||||
Token myToken = metadata.partitioner.getRandomToken();
|
||||
InetAddressAndPort myEndpoint = InetAddressAndPort.getByName("127.0.0.1");
|
||||
NodeId newNode = ClusterMetadataTestHelper.register(myEndpoint);
|
||||
ClusterMetadataTestHelper.JoinProcess join = ClusterMetadataTestHelper.lazyJoin(myEndpoint, myToken);
|
||||
join.prepareJoin();
|
||||
metadata = ClusterMetadata.current();
|
||||
BootstrapAndJoin joiningPlan = (BootstrapAndJoin) metadata.inProgressSequences.get(newNode);
|
||||
movements = joiningPlan.getMovementMaps(metadata);
|
||||
}
|
||||
|
||||
return new RangeStreamer(metadata,
|
||||
StreamOperation.BOOTSTRAP,
|
||||
true,
|
||||
DatabaseDescriptor.getNodeProximity(),
|
||||
new StreamStateStore(),
|
||||
mockFailureDetector,
|
||||
false,
|
||||
1,
|
||||
movements.left,
|
||||
movements.right);
|
||||
}
|
||||
|
||||
private boolean includesWraparound(Collection<Range<Token>> toFetch)
|
||||
{
|
||||
long minTokenCount = toFetch.stream()
|
||||
|
|
|
|||
Loading…
Reference in New Issue