diff --git a/CHANGES.txt b/CHANGES.txt index 3be443e292..b1518e8320 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.1.4 + * Memoize Cassandra verion and add a backoff interval for failed schema pulls (CASSANDRA-18902) * Fix StackOverflowError on ALTER after many previous schema changes (CASSANDRA-19166) * Fixed the inconsistency between distributedKeyspaces and distributedAndLocalKeyspaces (CASSANDRA-18747) * Internode legacy SSL storage port certificate is not hot reloaded on update (CASSANDRA-18681) diff --git a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java index 23f9598964..e51cdcbae0 100644 --- a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java +++ b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java @@ -453,6 +453,14 @@ public enum CassandraRelevantProperties return BOOLEAN_CONVERTER.convert(value); } + /** + * Clears the value set in the system property. + */ + public void clearValue() + { + System.clearProperty(key); + } + /** * Sets the value into system properties. * @param value to set diff --git a/src/java/org/apache/cassandra/schema/MigrationCoordinator.java b/src/java/org/apache/cassandra/schema/MigrationCoordinator.java index 61ef4c8cde..00ca3269f5 100644 --- a/src/java/org/apache/cassandra/schema/MigrationCoordinator.java +++ b/src/java/org/apache/cassandra/schema/MigrationCoordinator.java @@ -46,14 +46,12 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; import com.google.common.collect.ImmutableSet; import com.google.common.collect.Sets; -import org.apache.cassandra.utils.NoSpamLogger; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.concurrent.ExecutorPlus; import org.apache.cassandra.concurrent.FutureTask; import org.apache.cassandra.concurrent.ScheduledExecutors; -import org.apache.cassandra.config.CassandraRelevantProperties; import org.apache.cassandra.db.Mutation; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.gms.ApplicationState; @@ -68,6 +66,7 @@ import org.apache.cassandra.net.RequestCallback; import org.apache.cassandra.net.Verb; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.NoSpamLogger; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.Simulate; import org.apache.cassandra.utils.concurrent.Future; @@ -76,6 +75,7 @@ import org.apache.cassandra.utils.concurrent.WaitQueue; import static org.apache.cassandra.config.CassandraRelevantProperties.IGNORED_SCHEMA_CHECK_ENDPOINTS; import static org.apache.cassandra.config.CassandraRelevantProperties.IGNORED_SCHEMA_CHECK_VERSIONS; +import static org.apache.cassandra.config.CassandraRelevantProperties.MIGRATION_DELAY; import static org.apache.cassandra.config.CassandraRelevantProperties.SCHEMA_PULL_INTERVAL_MS; import static org.apache.cassandra.net.Verb.SCHEMA_PUSH_REQ; import static org.apache.cassandra.utils.Clock.Global.nanoTime; @@ -100,6 +100,7 @@ public class MigrationCoordinator private static final Logger logger = LoggerFactory.getLogger(MigrationCoordinator.class); private static final NoSpamLogger noSpamLogger = NoSpamLogger.getLogger(MigrationCoordinator.logger, 1, TimeUnit.MINUTES); private static final Future FINISHED_FUTURE = ImmediateFuture.success(null); + private static final long PULL_BACKOFF_INTERVAL_MS = 1000; // do not pull immediately if the previous pull failed private static LongSupplier getUptimeFn = () -> ManagementFactory.getRuntimeMXBean().getUptime(); @@ -109,7 +110,7 @@ public class MigrationCoordinator getUptimeFn = supplier; } - private static final int MIGRATION_DELAY_IN_MS = CassandraRelevantProperties.MIGRATION_DELAY.getInt(); + private static final int MIGRATION_DELAY_IN_MS = MIGRATION_DELAY.getInt(); public static final int MAX_OUTSTANDING_VERSION_REQUESTS = 3; private static ImmutableSet getIgnoredVersions() @@ -226,6 +227,7 @@ public class MigrationCoordinator private final Gossiper gossiper; private final Supplier schemaVersion; private final BiConsumer> schemaUpdateCallback; + private final Set lastPullFailures = new HashSet<>(); final ExecutorPlus executor; @@ -548,8 +550,16 @@ public class MigrationCoordinator if (shouldPullImmediately(endpoint, info.version)) { - logger.debug("Pulling {} immediately from {}", info, endpoint); - submitToMigrationIfNotShutdown(task); + if (lastPullFailures.contains(endpoint)) + { + logger.debug("Pulling {} immediately from {} with backoff interval = {}", info, endpoint, PULL_BACKOFF_INTERVAL_MS); + ScheduledExecutors.nonPeriodicTasks.schedule(() -> submitToMigrationIfNotShutdown(task), PULL_BACKOFF_INTERVAL_MS, TimeUnit.MILLISECONDS); + } + else + { + logger.debug("Pulling {} immediately from {}", info, endpoint); + submitToMigrationIfNotShutdown(task); + } } else { @@ -678,7 +688,14 @@ public class MigrationCoordinator private synchronized Future pullComplete(InetAddressAndPort endpoint, VersionInfo info, boolean wasSuccessful) { if (wasSuccessful) + { info.markReceived(); + lastPullFailures.remove(endpoint); + } + else + { + lastPullFailures.add(endpoint); + } info.outstandingRequests.remove(endpoint); info.requestQueue.add(endpoint); diff --git a/src/java/org/apache/cassandra/utils/FBUtilities.java b/src/java/org/apache/cassandra/utils/FBUtilities.java index 6d210ceb47..3d881ddf3f 100644 --- a/src/java/org/apache/cassandra/utils/FBUtilities.java +++ b/src/java/org/apache/cassandra/utils/FBUtilities.java @@ -37,6 +37,7 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.Future; import java.util.concurrent.TimeoutException; +import java.util.function.Supplier; import java.util.zip.CRC32; import java.util.zip.Checksum; @@ -45,6 +46,7 @@ import javax.annotation.Nullable; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Joiner; +import com.google.common.base.Suppliers; import com.fasterxml.jackson.databind.SerializationFeature; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; @@ -419,14 +421,11 @@ public class FBUtilities return previousReleaseVersionString; } - public static String getReleaseVersionString() - { + private static final Supplier loadedVersionString = Suppliers.memoize(() -> { try (InputStream in = FBUtilities.class.getClassLoader().getResourceAsStream("org/apache/cassandra/config/version.properties")) { if (in == null) - { - return System.getProperty("cassandra.releaseVersion", UNKNOWN_RELEASE_VERSION); - } + return null; Properties props = new Properties(); props.load(in); return props.getProperty("CassandraVersion"); @@ -437,6 +436,12 @@ public class FBUtilities logger.warn("Unable to load version.properties", e); return "debug version"; } + }); + + public static String getReleaseVersionString() + { + String v = loadedVersionString.get(); + return v != null ? v : System.getProperty("cassandra.releaseVersion", UNKNOWN_RELEASE_VERSION); } public static String getReleaseVersionMajor() diff --git a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java index bbb57f9c81..ab26654178 100644 --- a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java +++ b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java @@ -156,6 +156,7 @@ import org.apache.cassandra.utils.progress.jmx.JMXBroadcastExecutor; import static java.util.concurrent.TimeUnit.MINUTES; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; +import static org.apache.cassandra.config.CassandraRelevantProperties.RING_DELAY; import static org.apache.cassandra.distributed.api.Feature.BLANK_GOSSIP; import static org.apache.cassandra.distributed.api.Feature.GOSSIP; import static org.apache.cassandra.distributed.api.Feature.JMX; @@ -589,12 +590,20 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance // org.apache.cassandra.distributed.impl.AbstractCluster.startup sets the exception handler for the thread // so extract it to populate ExecutorFactory.Global ExecutorFactory.Global.tryUnsafeSet(new ExecutorFactory.Default(Thread.currentThread().getContextClassLoader(), null, Thread.getDefaultUncaughtExceptionHandler())); + + assert !FBUtilities.getReleaseVersionString().equals(FBUtilities.UNKNOWN_RELEASE_VERSION) : "Unknown version"; + assert !FBUtilities.getReleaseVersionString().isEmpty() : "Empty version"; + assert FBUtilities.getReleaseVersionString().contains(".") : "Invalid version: " + FBUtilities.getReleaseVersionString(); + if (config.has(GOSSIP)) { // TODO: hacky - System.setProperty("cassandra.ring_delay_ms", "15000"); - System.setProperty("cassandra.consistent.rangemovement", "false"); - System.setProperty("cassandra.consistent.simultaneousmoves.allow", "true"); + if (!RING_DELAY.isPresent()) + RING_DELAY.setLong(15000); + if (!System.getProperties().containsKey("cassandra.consistent.rangemovement")) + System.setProperty("cassandra.consistent.rangemovement", "false"); + if (!System.getProperties().containsKey("cassandra.consistent.simultaneousmoves.allow")) + System.setProperty("cassandra.consistent.simultaneousmoves.allow", "true"); } mkdirs(); diff --git a/test/distributed/org/apache/cassandra/distributed/test/MigrationCoordinatorTest.java b/test/distributed/org/apache/cassandra/distributed/test/MigrationCoordinatorTest.java index ca89b43d0c..f88792a9a7 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/MigrationCoordinatorTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/MigrationCoordinatorTest.java @@ -21,6 +21,7 @@ package org.apache.cassandra.distributed.test; import java.net.InetAddress; import java.util.UUID; +import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -32,21 +33,31 @@ import org.apache.cassandra.schema.Schema; import static org.apache.cassandra.config.CassandraRelevantProperties.IGNORED_SCHEMA_CHECK_ENDPOINTS; import static org.apache.cassandra.config.CassandraRelevantProperties.IGNORED_SCHEMA_CHECK_VERSIONS; +import static org.apache.cassandra.config.CassandraRelevantProperties.RING_DELAY; import static org.apache.cassandra.distributed.api.Feature.GOSSIP; import static org.apache.cassandra.distributed.api.Feature.NETWORK; public class MigrationCoordinatorTest extends TestBaseImpl { - @Before public void setUp() { - System.clearProperty("cassandra.replace_address"); - System.clearProperty("cassandra.consistent.rangemovement"); - - System.clearProperty(IGNORED_SCHEMA_CHECK_VERSIONS.getKey()); - System.clearProperty(IGNORED_SCHEMA_CHECK_VERSIONS.getKey()); + // make the test a bit faster + RING_DELAY.setLong(5000); + System.setProperty("cassandra.broadcast_interval_ms", "30000"); } + + @After + public void afterTestCleanup() + { + System.getProperties().remove("cassandra.replace_address"); + IGNORED_SCHEMA_CHECK_VERSIONS.clearValue(); + IGNORED_SCHEMA_CHECK_ENDPOINTS.clearValue(); + + RING_DELAY.clearValue(); + System.getProperties().remove("cassandra.broadcast_interval_ms"); + } + /** * We shouldn't wait on versions only available from a node being replaced * see CASSANDRA- @@ -89,7 +100,6 @@ public class MigrationCoordinatorTest extends TestBaseImpl IInstanceConfig config = cluster.newInstanceConfig(); config.set("auto_bootstrap", true); IGNORED_SCHEMA_CHECK_ENDPOINTS.setString(ignoredEndpoint.getHostAddress()); - System.setProperty("cassandra.consistent.rangemovement", "false"); cluster.bootstrap(config).startup(); } } @@ -116,7 +126,6 @@ public class MigrationCoordinatorTest extends TestBaseImpl IInstanceConfig config = cluster.newInstanceConfig(); config.set("auto_bootstrap", true); IGNORED_SCHEMA_CHECK_VERSIONS.setString(initialVersion.toString() + ',' + oldVersion); - System.setProperty("cassandra.consistent.rangemovement", "false"); cluster.bootstrap(config).startup(); } } diff --git a/test/unit/org/apache/cassandra/schema/MigrationCoordinatorTest.java b/test/unit/org/apache/cassandra/schema/MigrationCoordinatorTest.java index 3783320566..b4966af3a4 100644 --- a/test/unit/org/apache/cassandra/schema/MigrationCoordinatorTest.java +++ b/test/unit/org/apache/cassandra/schema/MigrationCoordinatorTest.java @@ -31,6 +31,7 @@ import java.util.UUID; import java.util.concurrent.Future; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; import com.google.common.collect.Iterables; import com.google.common.collect.Sets; @@ -40,6 +41,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.concurrent.ImmediateExecutor; +import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.Mutation; import org.apache.cassandra.gms.ApplicationState; @@ -56,6 +58,7 @@ import org.apache.cassandra.net.Verb; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.concurrent.WaitQueue; +import org.awaitility.Awaitility; import org.mockito.ArgumentCaptor; import org.mockito.ArgumentMatchers; import org.mockito.internal.creation.MockSettingsImpl; @@ -338,6 +341,8 @@ public class MigrationCoordinatorTest prev = next; Assert.assertFalse(wrapper.coordinator.awaitSchemaRequests(1)); + Awaitility.await().atMost(5, TimeUnit.SECONDS).pollDelay(100, TimeUnit.MILLISECONDS) + .until(() -> ScheduledExecutors.nonPeriodicTasks.getPendingTaskCount() == 0 && ScheduledExecutors.nonPeriodicTasks.getActiveTaskCount() == 0); Assert.assertEquals(2, wrapper.requests.size()); } logger.info("{} -> {}", EP1, ep1requests);