This is an automated email from the ASF dual-hosted git repository. aweisberg pushed a commit to branch cep-45-mutation-tracking in repository https://gitbox.apache.org/repos/asf/cassandra.git
commit a2e6a97edb26632a18be23de950a1d4ff3176e4c Author: Ariel Weisberg <[email protected]> AuthorDate: Wed Aug 12 11:54:39 2026 -0400 CEP-45: Log mutation tracking migration repair progress The repair completion callback named neither the repair nor the ranges, and nothing was logged when a keyspace finished migrating. All of it is at INFO, so a stalled migration can be debugged from a production log. - Session ids and the ranges repaired, remaining and already repaired, listed in full, once per repair job. - One line per keyspace when its last table finishes. - Why a repair did not contribute, checked before the epoch since an ineligible result carries Epoch.EMPTY. - Progress comes from the metadata commit returned, not current(), which races with other transformations. --- .../migration/KeyspaceMigrationInfo.java | 13 +++++ .../MutationTrackingMigrationRepairResult.java | 20 +++++-- .../migration/MutationTrackingRepairHandler.java | 63 ++++++++++++++++------ .../test/MutationTrackingMigrationTest.java | 18 +++++++ 4 files changed, 92 insertions(+), 22 deletions(-) diff --git a/src/java/org/apache/cassandra/service/replication/migration/KeyspaceMigrationInfo.java b/src/java/org/apache/cassandra/service/replication/migration/KeyspaceMigrationInfo.java index b0cb85b31a..65b08d0d78 100644 --- a/src/java/org/apache/cassandra/service/replication/migration/KeyspaceMigrationInfo.java +++ b/src/java/org/apache/cassandra/service/replication/migration/KeyspaceMigrationInfo.java @@ -218,6 +218,19 @@ public class KeyspaceMigrationInfo return ranges != null ? ranges : NormalizedRanges.empty(); } + /** The entire ring, which every table starts a migration needing to repair. */ + public static NormalizedRanges<Token> fullRing() + { + Token minimumToken = DatabaseDescriptor.getPartitioner().getMinimumToken(); + return NormalizedRanges.normalizedRanges(Collections.singleton(new Range<>(minimumToken, minimumToken))); + } + + /** Ranges that have finished migrating for a table: the full ring minus whatever is still pending. */ + public NormalizedRanges<Token> getMigratedRangesForTable(@Nonnull TableId tableId) + { + return fullRing().subtract(getPendingRangesForTable(tableId)); + } + /** * Check if token is in any pending range. * Used for routing decisions during migration. diff --git a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationRepairResult.java b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationRepairResult.java index da773035b2..e0ec934382 100644 --- a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationRepairResult.java +++ b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingMigrationRepairResult.java @@ -18,6 +18,8 @@ package org.apache.cassandra.service.replication.migration; +import javax.annotation.Nullable; + import org.apache.cassandra.tcm.Epoch; /** @@ -27,21 +29,29 @@ import org.apache.cassandra.tcm.Epoch; */ public class MutationTrackingMigrationRepairResult { - private static final MutationTrackingMigrationRepairResult INELIGIBLE = new MutationTrackingMigrationRepairResult(Epoch.EMPTY, false); + private static final MutationTrackingMigrationRepairResult DEAD_NODES_EXCLUDED = + new MutationTrackingMigrationRepairResult(Epoch.EMPTY, false, "dead nodes were excluded from the repair"); + private static final MutationTrackingMigrationRepairResult PREVIEW = + new MutationTrackingMigrationRepairResult(Epoch.EMPTY, false, "the repair was a preview"); public final Epoch minEpoch; public final boolean eligible; - private MutationTrackingMigrationRepairResult(Epoch minEpoch, boolean eligible) + /** Why this repair cannot contribute to migration, for logging. Null when eligible. */ + @Nullable + public final String ineligibleReason; + + private MutationTrackingMigrationRepairResult(Epoch minEpoch, boolean eligible, @Nullable String ineligibleReason) { this.minEpoch = minEpoch; this.eligible = eligible; + this.ineligibleReason = ineligibleReason; } public static MutationTrackingMigrationRepairResult fromRepair(Epoch minEpoch, boolean deadNodesExcluded, boolean isPreview) { - if (deadNodesExcluded) return INELIGIBLE; - if (isPreview) return INELIGIBLE; - return new MutationTrackingMigrationRepairResult(minEpoch, true); + if (deadNodesExcluded) return DEAD_NODES_EXCLUDED; + if (isPreview) return PREVIEW; + return new MutationTrackingMigrationRepairResult(minEpoch, true, null); } } diff --git a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingRepairHandler.java b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingRepairHandler.java index 2465268d8b..b55c362a30 100644 --- a/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingRepairHandler.java +++ b/src/java/org/apache/cassandra/service/replication/migration/MutationTrackingRepairHandler.java @@ -25,8 +25,10 @@ import com.google.common.util.concurrent.FutureCallback; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.dht.NormalizedRanges; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; +import org.apache.cassandra.repair.RepairJobDesc; import org.apache.cassandra.repair.RepairResult; import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.tcm.ClusterMetadata; @@ -49,9 +51,10 @@ public class MutationTrackingRepairHandler { try { - String keyspace = repairResult.desc.keyspace; - String tableName = repairResult.desc.columnFamily; - Collection<Range<Token>> repairedRanges = repairResult.desc.ranges; + RepairJobDesc desc = repairResult.desc; + String keyspace = desc.keyspace; + String tableName = desc.columnFamily; + Collection<Range<Token>> repairedRanges = desc.ranges; ClusterMetadata clusterMetadata = ClusterMetadata.current(); @@ -60,6 +63,8 @@ public class MutationTrackingRepairHandler if (migrationInfo == null) { + logger.info("Repair session {} (parent session {}) completed for {}.{} but the keyspace is not migrating, not advancing mutation tracking migration", + desc.sessionId, desc.parentSessionId, keyspace, tableName); return; } @@ -68,37 +73,61 @@ public class MutationTrackingRepairHandler if (tableMetadata == null) { - logger.warn("Repair completed for unknown table {}.{}, cannot advance migration", - keyspace, tableName); + logger.warn("Repair session {} (parent session {}) completed for unknown table {}.{}, cannot advance mutation tracking migration", + desc.sessionId, desc.parentSessionId, keyspace, tableName); return; } if (migrationInfo.getPendingRangesForTable(tableMetadata.id).isEmpty()) { - // Table already fully migrated + logger.info("Repair session {} (parent session {}) completed for {}.{} but the table has no ranges left to migrate, not advancing mutation tracking migration", + desc.sessionId, desc.parentSessionId, keyspace, tableName); return; } - // Epoch eligibility check: Only count repairs started after the migration started - if (repairResult.mutationTrackingMigrationRepairResult.minEpoch.isBefore(migrationInfo.startedAtEpoch)) + MutationTrackingMigrationRepairResult migrationRepairResult = repairResult.mutationTrackingMigrationRepairResult; + + // Before the epoch check: an ineligible result carries no epoch and would look stale + if (!migrationRepairResult.eligible) { - logger.debug("Repair completed for {}.{} but current epoch {} is before migration start epoch {}, ignoring", - keyspace, tableName, clusterMetadata.epoch, migrationInfo.startedAtEpoch); + logger.info("Repair session {} (parent session {}) completed for {}.{} but is ineligible to advance mutation tracking migration because {}", + desc.sessionId, desc.parentSessionId, keyspace, tableName, migrationRepairResult.ineligibleReason); return; } - if (!repairResult.mutationTrackingMigrationRepairResult.eligible) + // Epoch eligibility check: Only count repairs started after the migration started + if (migrationRepairResult.minEpoch.isBefore(migrationInfo.startedAtEpoch)) { - logger.debug("Repair completed for {}.{} but repair is ineligible for mutation tracking migration, ignoring", - keyspace, tableName); + logger.info("Repair session {} (parent session {}) completed for {}.{} but the repair started at epoch {}, before the migration started at epoch {}, not advancing mutation tracking migration", + desc.sessionId, desc.parentSessionId, keyspace, tableName, migrationRepairResult.minEpoch, migrationInfo.startedAtEpoch); return; } - logger.info("Repair completed for {}.{}, proposing migration advancement for {} ranges", - keyspace, tableName, repairedRanges.size()); - - ClusterMetadataService.instance().commit( + ClusterMetadata committed = ClusterMetadataService.instance().commit( new AdvanceMutationTrackingMigration(keyspace, tableMetadata.id, repairedRanges)); + + // Report from the metadata commit returned, not current(), which races with other epochs + KeyspaceMigrationInfo advanced = committed.mutationTrackingMigrationState.getKeyspaceInfo(keyspace); + boolean keyspaceComplete = advanced == null; + NormalizedRanges<Token> pending = keyspaceComplete ? NormalizedRanges.empty() + : advanced.getPendingRangesForTable(tableMetadata.id); + NormalizedRanges<Token> repaired = keyspaceComplete ? KeyspaceMigrationInfo.fullRing() + : advanced.getMigratedRangesForTable(tableMetadata.id); + + // INFO once per repair job, with the ranges listed in full rather than a prefix + logger.info("Repair session {} (parent session {}) advanced mutation tracking migration of {}.{} at epoch {}: " + + "contributed {} range(s) {}; {} range(s) remain to be repaired {}; {} range(s) already repaired {}; " + + "{} table(s) in the keyspace still migrating", + desc.sessionId, desc.parentSessionId, keyspace, tableName, committed.epoch, + repairedRanges.size(), repairedRanges, + pending.size(), pending, + repaired.size(), repaired, + keyspaceComplete ? 0 : advanced.pendingRangesPerTable.size()); + + // Only the advancement that empties the last table sees the keyspace disappear + if (keyspaceComplete) + logger.info("Mutation tracking migration completed for keyspace {} at epoch {}, every table has been fully repaired; final contribution from repair session {} (parent session {}) on table {}", + keyspace, committed.epoch, desc.sessionId, desc.parentSessionId, tableName); } catch (Exception e) { diff --git a/test/distributed/org/apache/cassandra/distributed/test/MutationTrackingMigrationTest.java b/test/distributed/org/apache/cassandra/distributed/test/MutationTrackingMigrationTest.java index c0df2f6779..0ffe1ede6e 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/MutationTrackingMigrationTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/MutationTrackingMigrationTest.java @@ -19,6 +19,8 @@ package org.apache.cassandra.distributed.test; import java.io.IOException; +import java.time.Duration; +import java.util.List; import java.util.concurrent.TimeoutException; import org.junit.BeforeClass; @@ -226,12 +228,28 @@ public class MutationTrackingMigrationTest extends TestBaseImpl assertTrue(journalEntriesAfterMigrationWrites > journalEntriesBeforeMigrationWrites); // complete migration + long logMark = SHARED_CLUSTER.get(1).logs().mark(); SHARED_CLUSTER.get(1).nodetoolResult("repair", testKeyspace, TEST_TABLE).asserts().success(); waitForEpochOf(SHARED_CLUSTER, 1); verifyKeyspaceState(testKeyspace, ExpectedKeyspaceState.TRACKED); + // Completion is logged once, and the contributing repair names itself and its ranges + List<String> completionLines = SHARED_CLUSTER.get(1).logs() + .watchFor(logMark, Duration.ofMinutes(1), "Mutation tracking migration completed for keyspace " + testKeyspace) + .getResult(); + assertEquals(1, completionLines.size()); + + List<String> advancementLines = SHARED_CLUSTER.get(1).logs() + .grep(logMark, "advanced mutation tracking migration of " + testKeyspace + '.' + TEST_TABLE) + .getResult(); + assertFalse(advancementLines.isEmpty()); + String advancement = advancementLines.get(0); + assertTrue(advancement, advancement.contains("parent session")); + assertTrue(advancement, advancement.contains("range(s) remain to be repaired")); + assertTrue(advancement, advancement.contains("range(s) already repaired")); + long journalEntriesBeforeTracked = countJournalEntries(); for (int i = 200; i < 210; i++) --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
