diff --git a/CHANGES.txt b/CHANGES.txt index 36b36b65be..c59faa6fe1 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -96,11 +96,8 @@ Merged from 4.0: 4.0.7 * Better handle null state in Gossip schema migration to avoid NPE (CASSANDRA-17864) * HintedHandoffAddRemoveNodesTest now accounts for the fact that StorageMetrics.totalHints is not updated synchronously w/ writes (CASSANDRA-16679) - - -4.0.7 * Avoid getting hanging repairs due to repair message timeouts (CASSANDRA-17613) - + * Fix resetting schema (CASSANDRA-17819) 4.0.6 * Prevent infinite loop in repair coordinator on FailSession (CASSANDRA-17834) diff --git a/NEWS.txt b/NEWS.txt index 9457c343bb..b64a06c35d 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -213,6 +213,9 @@ New features When the node is restarted with UUID based generation identifiers enabled, each newly created sstable will have a UUID based generation identifier and such files are not readable by previous Cassandra versions. In the future those new identifiers will become enabled by default. + - Resetting schema behavior has changed in 4.1 so that: 1) resetting schema is prohibited when there is no live node + where the schema could be fetched from, and 2) truncating local schema keyspace is postponed to the moment when + the node receives schema from some other node. Upgrading --------- diff --git a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java index 81f9a64418..49df7404e7 100644 --- a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java +++ b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java @@ -235,6 +235,8 @@ public enum CassandraRelevantProperties MEMTABLE_OVERHEAD_SIZE("cassandra.memtable.row_overhead_size", "-1"), MEMTABLE_OVERHEAD_COMPUTE_STEPS("cassandra.memtable_row_overhead_computation_step", "100000"), MIGRATION_DELAY("cassandra.migration_delay_ms", "60000"), + /** Defines how often schema definitions are pulled from the other nodes */ + SCHEMA_PULL_INTERVAL_MS("cassandra.schema_pull_interval_ms", "60000"), PAXOS_REPAIR_RETRY_TIMEOUT_IN_MS("cassandra.paxos_repair_retry_timeout_millis", "60000"), diff --git a/src/java/org/apache/cassandra/schema/DefaultSchemaUpdateHandler.java b/src/java/org/apache/cassandra/schema/DefaultSchemaUpdateHandler.java index 381b4e5ad9..0f0c3e9cf5 100644 --- a/src/java/org/apache/cassandra/schema/DefaultSchemaUpdateHandler.java +++ b/src/java/org/apache/cassandra/schema/DefaultSchemaUpdateHandler.java @@ -20,6 +20,7 @@ package org.apache.cassandra.schema; import java.time.Duration; import java.util.Collection; +import java.util.Collections; import java.util.Map; import java.util.Set; import java.util.UUID; @@ -46,6 +47,8 @@ import org.apache.cassandra.schema.SchemaTransformation.SchemaTransformationResu import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.Pair; +import org.apache.cassandra.utils.concurrent.AsyncPromise; +import org.apache.cassandra.utils.concurrent.Awaitable; import static org.apache.cassandra.schema.MigrationCoordinator.MAX_OUTSTANDING_VERSION_REQUESTS; @@ -60,6 +63,8 @@ public class DefaultSchemaUpdateHandler implements SchemaUpdateHandler, IEndpoin private final BiConsumer updateCallback; private volatile DistributedSchema schema = DistributedSchema.EMPTY; + private volatile AsyncPromise requestedReset; + private MigrationCoordinator createMigrationCoordinator(MessagingService messagingService) { return new MigrationCoordinator(messagingService, @@ -67,8 +72,8 @@ public class DefaultSchemaUpdateHandler implements SchemaUpdateHandler, IEndpoin ScheduledExecutors.scheduledTasks, MAX_OUTSTANDING_VERSION_REQUESTS, Gossiper.instance, - () -> schema.getVersion(), - (from, mutations) -> applyMutations(mutations)); + this::getSchemaVersionForCoordinator, + this::applyMutationsFromCoordinator); } public DefaultSchemaUpdateHandler(BiConsumer updateCallback) @@ -85,8 +90,23 @@ public class DefaultSchemaUpdateHandler implements SchemaUpdateHandler, IEndpoin this.updateCallback = updateCallback; this.migrationCoordinator = migrationCoordinator == null ? createMigrationCoordinator(messagingService) : migrationCoordinator; Gossiper.instance.register(this); - SchemaPushVerbHandler.instance.register(msg -> applyMutations(msg.payload)); - SchemaPullVerbHandler.instance.register(msg -> messagingService.send(msg.responseWith(getSchemaMutations()), msg.from())); + SchemaPushVerbHandler.instance.register(msg -> { + synchronized (this) + { + if (requestedReset == null) + applyMutations(msg.payload); + } + }); + SchemaPullVerbHandler.instance.register(msg -> { + try + { + messagingService.send(msg.responseWith(getSchemaMutations()), msg.from()); + } + catch (RuntimeException ex) + { + logger.error("Failed to send schema mutations to " + msg.from(), ex); + } + }); } public synchronized void start() @@ -198,6 +218,7 @@ public class DefaultSchemaUpdateHandler implements SchemaUpdateHandler, IEndpoin DistributedSchema after = new DistributedSchema(afterKeyspaces, version); SchemaTransformationResult update = new SchemaTransformationResult(before, after, diff); + logger.info("Applying schema change due to received mutations: {}", update); updateSchema(update, false); return update; } @@ -232,16 +253,23 @@ public class DefaultSchemaUpdateHandler implements SchemaUpdateHandler, IEndpoin private void updateSchema(SchemaTransformationResult update, boolean local) { - this.schema = update.after; - logger.debug("Schema updated: {}", update); - updateCallback.accept(update, true); - if (!local) + if (!update.diff.isEmpty()) { - migrationCoordinator.announce(update.after.getVersion()); + this.schema = update.after; + logger.debug("Schema updated: {}", update); + updateCallback.accept(update, true); + if (!local) + { + migrationCoordinator.announce(update.after.getVersion()); + } + } + else + { + logger.debug("Schema update is empty - skipping"); } } - private synchronized SchemaTransformationResult reload() + private synchronized void reload() { DistributedSchema before = this.schema; DistributedSchema after = new DistributedSchema(SchemaKeyspace.fetchNonSystemKeyspaces(), SchemaKeyspace.calculateSchemaDigest()); @@ -249,29 +277,76 @@ public class DefaultSchemaUpdateHandler implements SchemaUpdateHandler, IEndpoin SchemaTransformationResult update = new SchemaTransformationResult(before, after, diff); updateSchema(update, false); - return update; } @Override - public SchemaTransformationResult reset(boolean local) + public void reset(boolean local) { if (local) - return reload(); - - Collection mutations = migrationCoordinator.pullSchemaFromAnyNode().awaitThrowUncheckedOnInterrupt().getNow(); - return applyMutations(mutations); + { + reload(); + } + else + { + migrationCoordinator.reset(); + if (!migrationCoordinator.awaitSchemaRequests(CassandraRelevantProperties.MIGRATION_DELAY.getLong())) + { + logger.error("Timeout exceeded when waiting for schema from other nodes"); + } + } } + /** + * When clear is called the update handler will flag that the clear was requested. It means that migration + * coordinator will think that we have empty schema version and will apply whatever it receives from other nodes. + * When a first attempt to apply mutations from other node is called, it will first clear the schema and apply + * the mutations on a truncated table. The flag is then reset. + *

+ * This way the clear is postponed until we really fetch any schema we can use as a replacement. Otherwise, nothing + * will happen. We will simply reset the flag after the timeout and throw exceptions to the caller. + * + * @return + */ @Override - public synchronized void clear() + public Awaitable clear() { - SchemaKeyspace.truncate(); - this.schema = DistributedSchema.EMPTY; + synchronized (this) + { + if (requestedReset == null) + { + requestedReset = new AsyncPromise<>(); + migrationCoordinator.reset(); + } + return requestedReset; + } + } + + private UUID getSchemaVersionForCoordinator() + { + if (requestedReset != null) + return SchemaConstants.emptyVersion; + else + return schema.getVersion(); + } + + private synchronized void applyMutationsFromCoordinator(InetAddressAndPort from, Collection mutations) + { + if (requestedReset != null && !mutations.isEmpty()) + { + schema = DistributedSchema.EMPTY; + SchemaKeyspace.truncate(); + requestedReset.setSuccess(null); + requestedReset = null; + } + applyMutations(mutations); } private synchronized Collection getSchemaMutations() { - return SchemaKeyspace.convertSchemaToMutations(); + if (requestedReset != null) + return Collections.emptyList(); + else + return SchemaKeyspace.convertSchemaToMutations(); } public Map> getOutstandingSchemaVersions() diff --git a/src/java/org/apache/cassandra/schema/MigrationCoordinator.java b/src/java/org/apache/cassandra/schema/MigrationCoordinator.java index e57f24bbcd..fb592010c6 100644 --- a/src/java/org/apache/cassandra/schema/MigrationCoordinator.java +++ b/src/java/org/apache/cassandra/schema/MigrationCoordinator.java @@ -20,22 +20,22 @@ package org.apache.cassandra.schema; import java.lang.management.ManagementFactory; import java.net.UnknownHostException; +import java.time.Duration; import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Collection; -import java.util.Collections; import java.util.Deque; import java.util.HashMap; import java.util.HashSet; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.Optional; import java.util.Set; import java.util.UUID; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; -import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer; @@ -69,13 +69,13 @@ import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.Simulate; -import org.apache.cassandra.utils.concurrent.AsyncPromise; import org.apache.cassandra.utils.concurrent.Future; import org.apache.cassandra.utils.concurrent.ImmediateFuture; 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.SCHEMA_PULL_INTERVAL_MS; import static org.apache.cassandra.net.Verb.SCHEMA_PUSH_REQ; import static org.apache.cassandra.utils.Clock.Global.nanoTime; import static org.apache.cassandra.utils.Simulate.With.MONITORS; @@ -86,6 +86,10 @@ import static org.apache.cassandra.utils.concurrent.WaitQueue.newWaitQueue; * schema. It performs periodic checks and if there is a schema version mismatch between the current node and the other * node, it pulls the schema and applies the changes locally through the callback. * + * In particular the Migration Coordinator keeps track of all schema versions reported from each node in the cluster. + * As long as a certain version is advertised by some node, it is being tracked. As long as a version is tracked, + * the migration coordinator tries to fetch it by its periodic job. + * * It works in close cooperation with {@link DefaultSchemaUpdateHandler} which is responsible for maintaining local * schema metadata stored in {@link SchemaKeyspace}. */ @@ -150,12 +154,27 @@ public class MigrationCoordinator { final UUID version; + /** + * The set of endpoints containing this schema version + */ final Set endpoints = Sets.newConcurrentHashSet(); + /** + * The set of endpoints from which we are already fetching the schema + */ final Set outstandingRequests = Sets.newConcurrentHashSet(); + /** + * The queue of endpoints from which we are going to fetch the schema + */ final Deque requestQueue = new ArrayDeque<>(); + /** + * Threads waiting for schema synchronization are waiting until this object is signalled + */ private final WaitQueue waitQueue = newWaitQueue(); + /** + * Whether this schema version have been received + */ volatile boolean receivedSchema; VersionInfo(UUID version) @@ -181,6 +200,18 @@ public class MigrationCoordinator { return receivedSchema; } + + @Override + public String toString() + { + return "VersionInfo{" + + "version=" + version + + ", outstandingRequests=" + outstandingRequests + + ", requestQueue=" + requestQueue + + ", waitQueue.waiting=" + waitQueue.getWaiting() + + ", receivedSchema=" + receivedSchema + + '}'; + } } private final Map versionInfo = new HashMap<>(); @@ -221,18 +252,23 @@ public class MigrationCoordinator void start() { + logger.info("Starting migration coordinator and scheduling pulling schema versions every {}", Duration.ofMillis(SCHEMA_PULL_INTERVAL_MS.getLong())); announce(schemaVersion.get()); periodicPullTask.updateAndGet(curTask -> curTask == null - ? periodicCheckExecutor.scheduleWithFixedDelay(this::pullUnreceivedSchemaVersions, 1, 1, TimeUnit.MINUTES) + ? periodicCheckExecutor.scheduleWithFixedDelay(this::pullUnreceivedSchemaVersions, SCHEMA_PULL_INTERVAL_MS.getLong(), SCHEMA_PULL_INTERVAL_MS.getLong(), TimeUnit.MILLISECONDS) : curTask); } private synchronized void pullUnreceivedSchemaVersions() { + logger.debug("Pulling unreceived schema versions..."); for (VersionInfo info : versionInfo.values()) { if (info.wasReceived() || info.outstandingRequests.size() > 0) + { + logger.trace("Skipping pull of schema {} because it has been already recevied, or it is being received ({})", info.version, info); continue; + } maybePullSchema(info); } @@ -241,16 +277,25 @@ public class MigrationCoordinator private synchronized Future maybePullSchema(VersionInfo info) { if (info.endpoints.isEmpty() || info.wasReceived() || !shouldPullSchema(info.version)) + { + logger.trace("Not pulling schema {} because it was received, there is no endpoint to provide it, or we should not pull it ({})", info.version, info); return FINISHED_FUTURE; + } if (info.outstandingRequests.size() >= maxOutstandingVersionRequests) + { + logger.trace("Not pulling schema {} because the number of outstanding requests has been exceeded ({} >= {})", info.version, info.outstandingRequests.size(), maxOutstandingVersionRequests); return FINISHED_FUTURE; + } for (int i = 0, isize = info.requestQueue.size(); i < isize; i++) { InetAddressAndPort endpoint = info.requestQueue.remove(); if (!info.endpoints.contains(endpoint)) + { + logger.trace("Skipping request of schema {} from {} because the endpoint does not have that schema any longer", info.version, endpoint); continue; + } if (shouldPullFromEndpoint(endpoint) && info.outstandingRequests.add(endpoint)) { @@ -259,6 +304,7 @@ public class MigrationCoordinator else { // return to queue + logger.trace("Could not pull schema {} from {} - the request will be added back to the queue", info.version, endpoint); info.requestQueue.offer(endpoint); } } @@ -287,30 +333,33 @@ public class MigrationCoordinator UUID localSchemaVersion = schemaVersion.get(); if (localSchemaVersion == null) { - logger.debug("Not pulling schema for version {}, because local schama version is not known yet", version); + logger.debug("Not pulling schema {} because the local schama version is not known yet", version); return false; } if (localSchemaVersion.equals(version)) { - logger.debug("Not pulling schema for version {}, because schema versions match: " + - "local={}, remote={}", - version, - DistributedSchema.schemaVersionToString(localSchemaVersion), - DistributedSchema.schemaVersionToString(version)); + logger.debug("Not pulling schema {} because it is the same as the local schema", version); return false; } + return true; } private boolean shouldPullFromEndpoint(InetAddressAndPort endpoint) { if (endpoint.equals(FBUtilities.getBroadcastAddressAndPort())) + { + logger.trace("Not pulling schema from local endpoint"); return false; + } EndpointState state = gossiper.getEndpointStateForEndpoint(endpoint); if (state == null) + { + logger.trace("Not pulling schema from endpoint {} because its state is unknown", endpoint); return false; + } VersionedValue releaseVersionValue = state.getApplicationState(ApplicationState.RELEASE_VERSION); if (releaseVersionValue == null) @@ -373,22 +422,34 @@ public class MigrationCoordinator synchronized Future reportEndpointVersion(InetAddressAndPort endpoint, UUID version) { + logger.debug("Reported schema {} at endpoint {}", version, endpoint); if (ignoredEndpoints.contains(endpoint) || IGNORED_VERSIONS.contains(version)) { endpointVersions.remove(endpoint); removeEndpointFromVersion(endpoint, null); + logger.debug("Discarding endpoint {} or schema {} because either endpoint or schema version were marked as ignored", endpoint, version); return FINISHED_FUTURE; } UUID current = endpointVersions.put(endpoint, version); if (current != null && current.equals(version)) + { + logger.trace("Skipping report of schema {} from {} because we already know that", version, endpoint); return FINISHED_FUTURE; + } VersionInfo info = versionInfo.computeIfAbsent(version, VersionInfo::new); if (Objects.equals(schemaVersion.get(), version)) + { info.markReceived(); + logger.trace("Schema {} from {} has been marked as recevied because it is equal the local schema", version, endpoint); + } + else + { + info.requestQueue.addFirst(endpoint); + } info.endpoints.add(endpoint); - info.requestQueue.addFirst(endpoint); + logger.trace("Added endpoint {} to schema {}: {}", endpoint, info.version, info); // disassociate this endpoint from its (now) previous schema version removeEndpointFromVersion(endpoint, current); @@ -407,16 +468,70 @@ public class MigrationCoordinator return; info.endpoints.remove(endpoint); + logger.trace("Removed endpoint {} from schema {}: {}", endpoint, version, info); if (info.endpoints.isEmpty()) { info.waitQueue.signalAll(); versionInfo.remove(version); + logger.trace("Removed schema info: {}", info); } } + private void clearVersionsInfo() + { + Iterator> it = versionInfo.entrySet().iterator(); + while (it.hasNext()) + { + Map.Entry entry = it.next(); + it.remove(); + entry.getValue().waitQueue.signal(); + } + } + + private void reportCurrentSchemaVersionOnEndpoint(InetAddressAndPort endpoint) + { + if (FBUtilities.getBroadcastAddressAndPort().equals(endpoint)) + { + reportEndpointVersion(endpoint, schemaVersion.get()); + } + else + { + EndpointState state = gossiper.getEndpointStateForEndpoint(endpoint); + if (state != null) + { + UUID v = state.getSchemaVersion(); + if (v != null) + { + reportEndpointVersion(endpoint, v); + } + } + } + } + + /** + * Resets the migration coordinator by notifying all waiting threads and removing all the existing version info. + * Then, it is populated with the information about schema versions on different endpoints provided by Gossiper. + * Each version is marked as unreceived so the migration coordinator will start pulling schemas from other nodes. + */ + synchronized void reset() + { + logger.info("Resetting migration coordinator..."); + + // clear all the managed information + this.endpointVersions.clear(); + clearVersionsInfo(); + + // now report again the versions we are aware of + gossiper.getLiveMembers().forEach(this::reportCurrentSchemaVersionOnEndpoint); + } + synchronized void removeAndIgnoreEndpoint(InetAddressAndPort endpoint) { + logger.debug("Removing and ignoring endpoint {}", endpoint); Preconditions.checkArgument(endpoint != null); + // TODO The endpoint address is now ignored but when a node with the same address is added again later, + // there will be no way to include it in schema synchronization other than restarting each other node + // see https://issues.apache.org/jira/browse/CASSANDRA-17883 for details ignoredEndpoints.add(endpoint); Set versions = ImmutableSet.copyOf(versionInfo.keySet()); for (UUID version : versions) @@ -430,49 +545,19 @@ public class MigrationCoordinator FutureTask task = new FutureTask<>(() -> pullSchema(endpoint, new Callback(endpoint, info))); if (shouldPullImmediately(endpoint, info.version)) + { + logger.debug("Pulling {} immediately from {}", info, endpoint); submitToMigrationIfNotShutdown(task); + } else + { + logger.debug("Postponing pull of {} from {} for {}ms", info, endpoint, MIGRATION_DELAY_IN_MS); ScheduledExecutors.nonPeriodicTasks.schedule(() -> submitToMigrationIfNotShutdown(task), MIGRATION_DELAY_IN_MS, TimeUnit.MILLISECONDS); + } return task; } - private Future> pullSchemaFrom(InetAddressAndPort endpoint) - { - AsyncPromise> result = new AsyncPromise<>(); - return submitToMigrationIfNotShutdown(() -> pullSchema(endpoint, new RequestCallback>() - { - @Override - public void onResponse(Message> msg) - { - result.setSuccess(msg.payload); - } - - @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) - { - result.setFailure(new RuntimeException("Failed to get schema from " + from + ". The failure reason was: " + failureReason)); - } - - @Override - public boolean invokeOnFailure() - { - return true; - } - })).flatMap(ignored -> result); - } - - Future> pullSchemaFromAnyNode() - { - Optional endpoint = gossiper.getLiveMembers() - .stream() - .filter(this::shouldPullFromEndpoint) - .findFirst(); - - return endpoint.map(this::pullSchemaFrom).orElse(ImmediateFuture.success(Collections.emptyList())); - } - - void announce(UUID schemaVersion) { if (gossiper.isEnabled()) diff --git a/src/java/org/apache/cassandra/schema/OfflineSchemaUpdateHandler.java b/src/java/org/apache/cassandra/schema/OfflineSchemaUpdateHandler.java index 9d1020883b..16694b1f7c 100644 --- a/src/java/org/apache/cassandra/schema/OfflineSchemaUpdateHandler.java +++ b/src/java/org/apache/cassandra/schema/OfflineSchemaUpdateHandler.java @@ -27,6 +27,8 @@ import org.slf4j.LoggerFactory; import org.apache.cassandra.schema.SchemaTransformation.SchemaTransformationResult; import org.apache.cassandra.utils.ByteArrayUtil; +import org.apache.cassandra.utils.concurrent.Awaitable; +import org.apache.cassandra.utils.concurrent.ImmediateFuture; /** * Update handler which works only in memory. It does not load or save the schema anywhere. It is used in client mode @@ -77,17 +79,18 @@ public class OfflineSchemaUpdateHandler implements SchemaUpdateHandler } @Override - public SchemaTransformationResult reset(boolean local) + public void reset(boolean local) { if (!local) throw new UnsupportedOperationException(); - return apply(ignored -> SchemaKeyspace.fetchNonSystemKeyspaces(), local); + apply(ignored -> SchemaKeyspace.fetchNonSystemKeyspaces(), local); } @Override - public synchronized void clear() + public synchronized Awaitable clear() { this.schema = DistributedSchema.EMPTY; + return ImmediateFuture.success(true); } } diff --git a/src/java/org/apache/cassandra/schema/Schema.java b/src/java/org/apache/cassandra/schema/Schema.java index f23370c549..9869be4abf 100644 --- a/src/java/org/apache/cassandra/schema/Schema.java +++ b/src/java/org/apache/cassandra/schema/Schema.java @@ -19,6 +19,7 @@ package org.apache.cassandra.schema; import java.time.Duration; import java.util.*; +import java.util.concurrent.TimeUnit; import java.util.function.Consumer; import java.util.function.Supplier; @@ -34,6 +35,7 @@ import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.db.virtual.VirtualKeyspaceRegistry; import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.exceptions.InvalidRequestException; +import org.apache.cassandra.gms.Gossiper; import org.apache.cassandra.io.sstable.Descriptor; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.LocalStrategy; @@ -41,14 +43,16 @@ import org.apache.cassandra.schema.KeyspaceMetadata.KeyspaceDiff; import org.apache.cassandra.schema.Keyspaces.KeyspacesDiff; import org.apache.cassandra.schema.SchemaTransformation.SchemaTransformationResult; import org.apache.cassandra.service.PendingRangeCalculatorService; +import org.apache.cassandra.service.StorageService; +import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.concurrent.Awaitable; import org.apache.cassandra.utils.concurrent.LoadingMap; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import static java.lang.String.format; - import static com.google.common.collect.Iterables.size; +import static java.lang.String.format; import static org.apache.cassandra.config.DatabaseDescriptor.isDaemonInitialized; import static org.apache.cassandra.config.DatabaseDescriptor.isToolInitialized; @@ -611,18 +615,28 @@ public class Schema implements SchemaProvider * Clear all locally stored schema information and fetch schema from another node. * Called by user (via JMX) who wants to get rid of schema disagreement. */ - public synchronized void resetLocalSchema() + public void resetLocalSchema() { logger.debug("Clearing local schema..."); - updateHandler.clear(); - logger.debug("Clearing local schema keyspace instances..."); - distributedKeyspaces.forEach(this::unload); - updateVersion(SchemaConstants.emptyVersion); + if (Gossiper.instance.getLiveMembers().stream().allMatch(ep -> FBUtilities.getBroadcastAddressAndPort().equals(ep))) + throw new InvalidRequestException("Cannot reset local schema when there are no other live nodes"); + + Awaitable clearCompletion = updateHandler.clear(); + try + { + if (!clearCompletion.await(StorageService.SCHEMA_DELAY_MILLIS, TimeUnit.MILLISECONDS)) + { + throw new RuntimeException("Schema reset failed - no schema received from other nodes"); + } + } + catch (InterruptedException e) + { + Thread.currentThread().interrupt(); + throw new RuntimeException("Failed to reset schema - the thread has been interrupted"); + } SchemaDiagnostics.schemaCleared(this); - - updateHandler.reset(false); - logger.info("Local schema reset is complete."); + logger.info("Local schema reset completed"); } private void merge(KeyspacesDiff diff, boolean removeData) diff --git a/src/java/org/apache/cassandra/schema/SchemaUpdateHandler.java b/src/java/org/apache/cassandra/schema/SchemaUpdateHandler.java index c429f9b5d5..f6711f3219 100644 --- a/src/java/org/apache/cassandra/schema/SchemaUpdateHandler.java +++ b/src/java/org/apache/cassandra/schema/SchemaUpdateHandler.java @@ -21,6 +21,7 @@ package org.apache.cassandra.schema; import java.time.Duration; import org.apache.cassandra.schema.SchemaTransformation.SchemaTransformationResult; +import org.apache.cassandra.utils.concurrent.Awaitable; /** * Schema update handler is responsible for maintaining the shared schema and synchronizing it with other nodes in @@ -63,14 +64,15 @@ public interface SchemaUpdateHandler * refreshed, the callbacks provided in the factory method are executed, and the updated schema version is announced. * * @param local whether we should reset with locally stored schema or fetch the schema from other nodes - * @return transformation result */ - SchemaTransformationResult reset(boolean local); + void reset(boolean local); /** - * Clears the locally stored schema entirely. After this operation the schema is equal to {@link DistributedSchema#EMPTY}. - * The method does not execute any callback. It is indended to reinitialize the schema later using the method - * {@link #reset(boolean)}. + * Marks the local schema to be cleared and refreshed. Since calling this method, the update handler tries to obtain + * a fresh schema definition from a remote source. Once the schema definition is received, the local schema is + * replaced (instead of being merged which usually happens when the update is received). + *

+ * The returned awaitable is fulfilled when the schema is received and applied. */ - void clear(); + Awaitable clear(); } diff --git a/test/distributed/org/apache/cassandra/distributed/test/SchemaTest.java b/test/distributed/org/apache/cassandra/distributed/test/SchemaTest.java index 7b03105063..142b66b581 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/SchemaTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/SchemaTest.java @@ -18,20 +18,35 @@ package org.apache.cassandra.distributed.test; -import java.time.Duration; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import org.junit.Test; +import org.apache.cassandra.config.CassandraRelevantProperties; +import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.distributed.Cluster; import org.apache.cassandra.distributed.api.ConsistencyLevel; import org.apache.cassandra.distributed.api.Feature; +import org.apache.cassandra.distributed.api.IInvokableInstance; +import org.apache.cassandra.distributed.api.IIsolatedExecutor.SerializableCallable; +import org.apache.cassandra.gms.Gossiper; import org.apache.cassandra.schema.Schema; +import org.apache.cassandra.schema.SchemaConstants; +import org.assertj.core.api.Assertions; import org.awaitility.Awaitility; +import org.awaitility.core.ConditionFactory; +import static java.time.Duration.ofSeconds; +import static org.apache.cassandra.utils.FBUtilities.getBroadcastAddressAndPort; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; public class SchemaTest extends TestBaseImpl { + public static final String TABLE_ONE = "tbl_one"; + public static final String TABLE_TWO = "tbl_two"; + @Test public void readRepair() throws Throwable { @@ -84,7 +99,7 @@ public class SchemaTest extends TestBaseImpl Throwable cause = e; while (cause != null) { - if (cause.getMessage() != null && cause.getMessage().contains("Unknown column "+name+" during deserialization")) + if (cause.getMessage() != null && cause.getMessage().contains("Unknown column " + name + " during deserialization")) causeIsUnknownColumn = true; cause = cause.getCause(); } @@ -92,32 +107,94 @@ public class SchemaTest extends TestBaseImpl } } + /** + * The purpose of this test is to verify manual schema reset functinality. + *

+ * There is a 2-node cluster and a TABLE_ONE created. The schema version is agreed on both nodes. Then the 2nd node + * is shutdown. We introduce a disagreement by dropping TABLE_ONE and creating TABLE_TWO on the 1st node. Therefore, + * the 1st node has a newer schema version with TABLE_TWO, while the shutdown 2nd node has older schema version with + * TABLE_ONE. + *

+ * At this point, if we just started the 2nd node, it would sync its schema by getting fresh mutations from the 1st + * node which would result in both nodes having only the definition of TABLE_TWO. + *

+ * However, before starting the 2nd node the schema is reset on the 1st node, so the 1st node will discard its local + * schema whenever it manages to fetch a schema definition from some other node (the 2nd node in this case). + * It is expected to end up with both nodes having only the definition of TABLE_ONE. + *

+ * In the second phase of the test we simply break the schema on the 1st node and call reset to fetch the schema + * definition it from the 2nd node. + */ @Test public void schemaReset() throws Throwable { + CassandraRelevantProperties.MIGRATION_DELAY.setLong(10000); + CassandraRelevantProperties.SCHEMA_PULL_INTERVAL_MS.setLong(10000); try (Cluster cluster = init(Cluster.build(2).withConfig(cfg -> cfg.with(Feature.GOSSIP, Feature.NETWORK)).start())) { - cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (pk INT PRIMARY KEY, v TEXT)"); - - assertTrue(cluster.get(1).callOnInstance(() -> Schema.instance.getTableMetadata(KEYSPACE, "tbl") != null)); - assertTrue(cluster.get(2).callOnInstance(() -> Schema.instance.getTableMetadata(KEYSPACE, "tbl") != null)); + // create TABLE_ONE and make sure it is propagated + cluster.schemaChange(String.format("CREATE TABLE %s.%s (pk INT PRIMARY KEY, v TEXT)", KEYSPACE, TABLE_ONE)); + assertTrue(checkTablesPropagated(cluster.get(1), true, false)); + assertTrue(checkTablesPropagated(cluster.get(2), true, false)); + // shutdown the 2nd node and make sure that the 1st does not see it any longer as alive cluster.get(2).shutdown().get(); + await(30).until(() -> cluster.get(1).callOnInstance(() -> { + return Gossiper.instance.getLiveMembers() + .stream() + .allMatch(e -> e.equals(getBroadcastAddressAndPort())); + })); - // when schema is removed and there is no other node to fetch it from, node 1 should be left with clean schema - cluster.get(1).runOnInstance(() -> Schema.instance.resetLocalSchema()); - assertTrue(cluster.get(1).callOnInstance(() -> Schema.instance.getTableMetadata(KEYSPACE, "tbl") == null)); + // when there is no node to fetch the schema from, reset local schema should immediately fail + Assertions.assertThatExceptionOfType(RuntimeException.class).isThrownBy(() -> { + cluster.get(1).runOnInstance(() -> Schema.instance.resetLocalSchema()); + }).withMessageContaining("Cannot reset local schema when there are no other live nodes"); - // when the other node is started, schema should be back in sync + // now, let's make a disagreement, the shutdown node 2 has a definition of TABLE_ONE, while the running + // node 1 will have a definition of TABLE_TWO + cluster.coordinator(1).execute(String.format("DROP TABLE %s.%s", KEYSPACE, TABLE_ONE), ConsistencyLevel.ONE); + cluster.coordinator(1).execute(String.format("CREATE TABLE %s.%s (pk INT PRIMARY KEY, v TEXT)", KEYSPACE, TABLE_TWO), ConsistencyLevel.ONE); + await(30).until(() -> checkTablesPropagated(cluster.get(1), false, true)); + + // Schema.resetLocalSchema is guarded by some conditions which would not let us reset schema if there is no + // live node in the cluster, therefore we simply call SchemaUpdateHandler.clear (this is the only real thing + // being done by Schema.resetLocalSchema under the hood) + SerializableCallable clear = () -> Schema.instance.updateHandler.clear().awaitUninterruptibly(1, TimeUnit.MINUTES); + Future clear1 = cluster.get(1).asyncCallsOnInstance(clear).call(); + assertFalse(clear1.isDone()); + + // when the 2nd node is started, schema should be back in sync cluster.get(2).startup(); - Awaitility.waitAtMost(Duration.ofMinutes(1)) - .pollDelay(Duration.ofSeconds(1)) - .until(() -> cluster.get(1).callOnInstance(() -> Schema.instance.getTableMetadata(KEYSPACE, "tbl") != null)); + await(30).until(() -> clear1.isDone() && clear1.get()); - // when schema is removed and there is a node to fetch it from, node 1 should immediatelly restore the schema + // this proves that reset schema works on the 1st node - the most recent change should be discarded because + // it receives the schema from the 2nd node and applies it on empty schema + await(60).until(() -> checkTablesPropagated(cluster.get(1), true, false)); + + // now let's break schema locally and let it be reset + cluster.get(1).runOnInstance(() -> Schema.instance.getLocalKeyspaces() + .get(SchemaConstants.SCHEMA_KEYSPACE_NAME) + .get().tables.forEach(t -> ColumnFamilyStore.getIfExists(t.keyspace, t.name).truncateBlockingWithoutSnapshot())); + + // when schema is removed and there is a node to fetch it from, the 1st node should immediately restore it cluster.get(1).runOnInstance(() -> Schema.instance.resetLocalSchema()); - assertTrue(cluster.get(1).callOnInstance(() -> Schema.instance.getTableMetadata(KEYSPACE, "tbl") != null)); + // note that we should not wait for this to be true because resetLocalSchema is blocking + // and after successfully completing it, the schema should be already back in sync + assertTrue(checkTablesPropagated(cluster.get(1), true, false)); + assertTrue(checkTablesPropagated(cluster.get(2), true, false)); } } + private static ConditionFactory await(int seconds) + { + return Awaitility.await().atMost(ofSeconds(seconds)).pollDelay(ofSeconds(1)); + } + + private static boolean checkTablesPropagated(IInvokableInstance instance, boolean one, boolean two) + { + return instance.callOnInstance(() -> { + return (Schema.instance.getTableMetadata(KEYSPACE, TABLE_ONE) != null ^ !one) + && (Schema.instance.getTableMetadata(KEYSPACE, TABLE_TWO) != null ^ !two); + }); + } } diff --git a/test/distributed/org/apache/cassandra/distributed/test/jmx/JMXGetterCheckTest.java b/test/distributed/org/apache/cassandra/distributed/test/jmx/JMXGetterCheckTest.java index 217ec4cf7a..ef0f1ecae0 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/jmx/JMXGetterCheckTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/jmx/JMXGetterCheckTest.java @@ -53,7 +53,8 @@ public class JMXGetterCheckTest extends TestBaseImpl private static final Set IGNORE_OPERATIONS = ImmutableSet.of( "org.apache.cassandra.db:type=StorageService:stopDaemon", // halts the instance, which then causes the JVM to exit "org.apache.cassandra.db:type=StorageService:drain", // don't drain, it stops things which can cause other APIs to be unstable as we are in a stopped state - "org.apache.cassandra.db:type=StorageService:stopGossiping" // if we stop gossip this can causes other issues, so avoid + "org.apache.cassandra.db:type=StorageService:stopGossiping", // if we stop gossip this can cause other issues, so avoid + "org.apache.cassandra.db:type=StorageService:resetLocalSchema" // this will fail when there are no other nodes which can serve schema ); @Test diff --git a/test/unit/org/apache/cassandra/schema/MigrationCoordinatorTest.java b/test/unit/org/apache/cassandra/schema/MigrationCoordinatorTest.java index b9980414a7..3783320566 100644 --- a/test/unit/org/apache/cassandra/schema/MigrationCoordinatorTest.java +++ b/test/unit/org/apache/cassandra/schema/MigrationCoordinatorTest.java @@ -42,7 +42,6 @@ import org.slf4j.LoggerFactory; import org.apache.cassandra.concurrent.ImmediateExecutor; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.Mutation; -import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.gms.ApplicationState; import org.apache.cassandra.gms.EndpointState; import org.apache.cassandra.gms.Gossiper; @@ -63,7 +62,6 @@ import org.mockito.internal.creation.MockSettingsImpl; import static com.google.common.util.concurrent.Futures.getUnchecked; import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.assertThatExceptionOfType; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.Mockito.doAnswer; @@ -398,25 +396,27 @@ public class MigrationCoordinatorTest } @Test - public void pullSchemaFromAnyNode() throws UnknownHostException + public void reset() throws UnknownHostException { Collection mutations = Arrays.asList(mock(Mutation.class)); Wrapper wrapper = new Wrapper(); - - // no live nodes - when(wrapper.gossiper.getLiveMembers()).thenReturn(Collections.emptySet()); - Collection result = wrapper.coordinator.pullSchemaFromAnyNode().syncThrowUncheckedOnInterrupt().getNow(); - assertThat(result).isEmpty(); + wrapper.localSchemaVersion = SchemaConstants.emptyVersion; EndpointState invalidVersionState = mock(EndpointState.class); when(invalidVersionState.getApplicationState(ApplicationState.RELEASE_VERSION)).thenReturn(VersionedValue.unsafeMakeVersionedValue("3.0", 0)); + when(invalidVersionState.getSchemaVersion()).thenReturn(V1); EndpointState validVersionState = mock(EndpointState.class); when(validVersionState.getApplicationState(ApplicationState.RELEASE_VERSION)).thenReturn(VersionedValue.unsafeMakeVersionedValue(FBUtilities.getReleaseVersionString(), 0)); + when(validVersionState.getSchemaVersion()).thenReturn(V2); + + EndpointState localVersionState = mock(EndpointState.class); + when(localVersionState.getApplicationState(ApplicationState.RELEASE_VERSION)).thenReturn(VersionedValue.unsafeMakeVersionedValue(FBUtilities.getReleaseVersionString(), 0)); + when(localVersionState.getSchemaVersion()).thenReturn(SchemaConstants.emptyVersion); // some nodes - InetAddressAndPort thisNode = wrapper.configureMocksForEndpoint(FBUtilities.getBroadcastAddressAndPort(), validVersionState, MessagingService.current_version, false); + InetAddressAndPort thisNode = wrapper.configureMocksForEndpoint(FBUtilities.getBroadcastAddressAndPort(), localVersionState, MessagingService.current_version, false); InetAddressAndPort noStateNode = wrapper.configureMocksForEndpoint("10.0.0.1:8000", null, MessagingService.current_version, false); InetAddressAndPort diffMajorVersionNode = wrapper.configureMocksForEndpoint("10.0.0.2:8000", invalidVersionState, MessagingService.current_version, false); InetAddressAndPort unkonwnNode = wrapper.configureMocksForEndpoint("10.0.0.2:8000", validVersionState, null, false); @@ -436,22 +436,8 @@ public class MigrationCoordinatorTest callback.onResponse(Message.remoteResponse(regularNode1, Verb.SCHEMA_PULL_RSP, mutations)); return null; }).when(wrapper.messagingService).sendWithCallback(any(Message.class), any(InetAddressAndPort.class), any(RequestCallback.class)); - result = wrapper.coordinator.pullSchemaFromAnyNode().syncThrowUncheckedOnInterrupt().getNow(); - assertThat(result).isEqualTo(mutations); - - // failures - doAnswer(a -> { - Message msg = a.getArgument(0, Message.class); - InetAddressAndPort endpoint = a.getArgument(1, InetAddressAndPort.class); - RequestCallback callback = a.getArgument(2, RequestCallback.class); - - assertThat(msg.verb()).isEqualTo(Verb.SCHEMA_PULL_REQ); - assertThat(endpoint).isEqualTo(regularNode1); - callback.onFailure(regularNode1, RequestFailureReason.UNKNOWN); - return null; - }).when(wrapper.messagingService).sendWithCallback(any(Message.class), any(InetAddressAndPort.class), any(RequestCallback.class)); - assertThatExceptionOfType(RuntimeException.class).isThrownBy(() -> wrapper.coordinator.pullSchemaFromAnyNode().syncThrowUncheckedOnInterrupt().getNow()) - .withMessageContaining("Failed to get schema from"); - + wrapper.coordinator.reset(); + assertThat(wrapper.mergedSchemasFrom).anyMatch(ep -> regularNode1.equals(ep) || regularNode2.equals(ep)); + assertThat(wrapper.mergedSchemasFrom).hasSize(1); } } \ No newline at end of file