From 53eb2fdac886fb559ce36649d97c3a87a0423e58 Mon Sep 17 00:00:00 2001 From: Alex Petrov Date: Mon, 14 Jul 2025 19:18:27 +0200 Subject: [PATCH] Improve GC testing; fix durability startup order Patch by Alex Petrov; reviewed by Benedict Elliott Smith for CASSANDRA-20767 --- modules/accord | 2 +- .../apache/cassandra/journal/Compactor.java | 2 +- .../AbstractAccordSegmentCompactor.java | 1 - .../service/accord/AccordDataStore.java | 7 ++ .../service/accord/AccordJournal.java | 10 +-- .../service/accord/AccordService.java | 5 +- .../fuzz/topology/JournalGCTest.java | 76 +++++++++++-------- 7 files changed, 64 insertions(+), 39 deletions(-) diff --git a/modules/accord b/modules/accord index 51f36f788b..e66a30cd94 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 51f36f788b8095c0da9efaf1ea91bb9a7dd31181 +Subproject commit e66a30cd944e902ad2749b5182e1f0d56903303f diff --git a/src/java/org/apache/cassandra/journal/Compactor.java b/src/java/org/apache/cassandra/journal/Compactor.java index 6525df5d05..330ecea0a9 100644 --- a/src/java/org/apache/cassandra/journal/Compactor.java +++ b/src/java/org/apache/cassandra/journal/Compactor.java @@ -70,7 +70,7 @@ public final class Compactor implements Runnable, Shutdownable { Set> toCompact = new HashSet<>(); journal.segments().selectStatic(toCompact); - if (toCompact.size() < 2) + if (toCompact.isEmpty()) return; try diff --git a/src/java/org/apache/cassandra/service/accord/AbstractAccordSegmentCompactor.java b/src/java/org/apache/cassandra/service/accord/AbstractAccordSegmentCompactor.java index c26ae3f679..0c26521f8b 100644 --- a/src/java/org/apache/cassandra/service/accord/AbstractAccordSegmentCompactor.java +++ b/src/java/org/apache/cassandra/service/accord/AbstractAccordSegmentCompactor.java @@ -96,7 +96,6 @@ public abstract class AbstractAccordSegmentCompactor implements SegmentCompac @Override public Collection> compact(Collection> segments) { - Invariants.require(segments.size() >= 2, () -> String.format("Can only compact 2 or more segments, but got %d", segments.size())); logger.info("Compacting {} static segments: {}", segments.size(), segments); // TODO (expected): this will be a large over-estimate. should make segments an sstable format and include cardinality estimation diff --git a/src/java/org/apache/cassandra/service/accord/AccordDataStore.java b/src/java/org/apache/cassandra/service/accord/AccordDataStore.java index 507f0f91cd..83e92321fb 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordDataStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordDataStore.java @@ -69,6 +69,13 @@ public class AccordDataStore implements DataStore while (true) { Memtable memtable = cfs.getCurrentMemtable(); + // If RX came when after a quiet period or if it raced with a previous memtable flush + if (memtable.isClean()) + { + AccordDurableOnFlush.notify(cfs.metadata(), commandStore, reportOnSuccess); + break; + } + AccordDurableOnFlush onFlush = memtable.ensureFlushListener(FlushListenerKey.KEY, AccordDurableOnFlush::new); if (onFlush != null && onFlush.add(commandStore.id(), reportOnSuccess)) break; diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournal.java b/src/java/org/apache/cassandra/service/accord/AccordJournal.java index 3e4ec84e07..08819934d0 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournal.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournal.java @@ -36,7 +36,6 @@ import java.util.function.Supplier; import javax.annotation.Nullable; import com.google.common.annotations.VisibleForTesting; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -376,17 +375,18 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier journal.onDurable(pointer, onFlush); } + private static final JournalKey DURABLE_BEFORE_KEY = new JournalKey(TxnId.NONE, JournalKey.Type.DURABLE_BEFORE, 0); + @Override public PersistentField.Persister durableBeforePersister() { return new PersistentField.Persister<>() { @Override - public AsyncResult persist(DurableBefore addDurableBefore, DurableBefore newDurableBefore) + public AsyncResult persist(DurableBefore addValue, DurableBefore newValue) { AsyncResult.Settable result = AsyncResults.settable(); - JournalKey key = new JournalKey(TxnId.NONE, JournalKey.Type.DURABLE_BEFORE, 0); - RecordPointer pointer = appendInternal(key, addDurableBefore); + RecordPointer pointer = appendInternal(DURABLE_BEFORE_KEY, addValue); // TODO (required): what happens on failure? journal.onDurable(pointer, () -> result.setSuccess(null)); return result; @@ -395,7 +395,7 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier @Override public DurableBefore load() { - DurableBeforeAccumulator accumulator = readAll(new JournalKey(TxnId.NONE, JournalKey.Type.DURABLE_BEFORE, 0)); + DurableBeforeAccumulator accumulator = readAll(DURABLE_BEFORE_KEY); return accumulator.get(); } }; diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index 3d20f18677..4432beb0aa 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -248,6 +248,10 @@ public class AccordService implements IAccordService, Shutdownable instance = as; replayJournal(as); + + // Only enable durability scheduling _after_ we have fully replayed journal + as.configService.registerListener(as.node.durability()); + as.node.durability().start(); } @VisibleForTesting @@ -456,7 +460,6 @@ public class AccordService implements IAccordService, Shutdownable Ints.checkedCast(getAccordShardDurabilityMaxSplits()), Ints.checkedCast(getAccordShardDurabilityCycle(SECONDS)), SECONDS); node.durability().global().setGlobalCycleTime(Ints.checkedCast(getAccordGlobalDurabilityCycle(SECONDS)), SECONDS); - node.durability().start(); state = State.STARTED; } diff --git a/test/distributed/org/apache/cassandra/fuzz/topology/JournalGCTest.java b/test/distributed/org/apache/cassandra/fuzz/topology/JournalGCTest.java index e685f3a4f9..03c3152808 100644 --- a/test/distributed/org/apache/cassandra/fuzz/topology/JournalGCTest.java +++ b/test/distributed/org/apache/cassandra/fuzz/topology/JournalGCTest.java @@ -18,9 +18,16 @@ package org.apache.cassandra.fuzz.topology; +import java.util.concurrent.Callable; +import java.util.concurrent.atomic.AtomicInteger; + +import org.junit.Assert; +import org.junit.Test; + import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.distributed.Cluster; import org.apache.cassandra.distributed.api.ConsistencyLevel; +import org.apache.cassandra.distributed.shared.ClusterUtils; import org.apache.cassandra.distributed.test.log.FuzzTestBase; import org.apache.cassandra.harry.SchemaSpec; import org.apache.cassandra.harry.dsl.HistoryBuilder; @@ -34,10 +41,6 @@ import org.apache.cassandra.service.accord.AccordKeyspace; import org.apache.cassandra.service.accord.AccordService; import org.apache.cassandra.service.accord.JournalKey; import org.apache.cassandra.service.consensus.TransactionalMode; -import org.junit.Assert; -import org.junit.Test; - -import java.util.concurrent.atomic.AtomicInteger; import static org.apache.cassandra.harry.checker.TestHelper.withRandom; @@ -49,10 +52,13 @@ public class JournalGCTest extends FuzzTestBase public void journalGCTest() throws Throwable { try (Cluster cluster = init(builder().withNodes(1) - .withConfig(cfg -> cfg.set("accord.gc_delay", "1s") - .set("accord.shard_durability_target_splits", "1") - .set("accord.shard_durability_cycle", "1s") - .set("accord.global_durability_cycle", "1s")) + .withConfig(cfg -> cfg.set("write_request_timeout", "2s") + .set("accord.expire_syncpoint", "1s*attempts<=300s") + .set("accord.retry_syncpoint", "1s*attempts") + .set("accord.shard_durability_target_splits", "5") + .set("accord.shard_durability_max_splits", "10") + .set("accord.shard_durability_cycle", "1s") + .set("accord.global_durability_cycle", "1s")) .start())) { withRandom(rng -> { @@ -74,9 +80,23 @@ public class JournalGCTest extends FuzzTestBase .pageSizeSelector(p -> InJvmDTestVisitExecutor.PageSizeSelector.NO_PAGING) .build(schema, hb, cluster)); - for (int pk = 0; pk < 500; pk++) { - for (int i = 0; i < 500; i++) + for (int pk = 0; pk <= 500; pk++) { + for (int i = 0; i < 100; i++) history.insert(pk); + + if (pk > 0 && pk % 100 == 0) + { + cluster.get(1).runOnInstance(() -> { + ((AccordService) AccordService.instance()).journal().closeCurrentSegmentForTestingIfNonEmpty(); + ((AccordService) AccordService.instance()).journal().compactor().run(); + }); + } + + if (pk > 0 && pk % 200 == 0) + { + ClusterUtils.stopUnchecked(cluster.get(1)); + cluster.get(1).startup(); + } } cluster.get(1).runOnInstance(() -> { @@ -84,34 +104,30 @@ public class JournalGCTest extends FuzzTestBase ((AccordService) AccordService.instance()).journal().compactor().run(); }); - int before = cluster.get(1).callOnInstance(() -> { + Callable countDiffs = cluster.get(1).callsOnInstance(() -> { AtomicInteger a = new AtomicInteger(); ((AccordService) AccordService.instance()).journal().forEach((v) -> { - if (v.type == JournalKey.Type.COMMAND_DIFF) + if (v.type == JournalKey.Type.COMMAND_DIFF && + // Do not count syncpoints + !v.id.isSyncPoint()) a.incrementAndGet(); }); return a.get(); }); - Thread.sleep(10_000); - cluster.get(1).runOnInstance(() -> { - Keyspace.open(SchemaConstants.ACCORD_KEYSPACE_NAME).getColumnFamilyStore(AccordKeyspace.JOURNAL).forceMajorCompaction(); - }); - - cluster.get(1).forceCompact("system_accord", "journal"); - - int after = cluster.get(1).callOnInstance(() -> { - AtomicInteger a = new AtomicInteger(); - ((AccordService) AccordService.instance()).journal().forEach((v) -> { - if (v.type == JournalKey.Type.COMMAND_DIFF) - a.incrementAndGet(); + int after =-1; + for (int i = 0; i < 60; i++) + { + cluster.get(1).runOnInstance(() -> { + Keyspace.open(SchemaConstants.ACCORD_KEYSPACE_NAME).getColumnFamilyStore(AccordKeyspace.JOURNAL).forceMajorCompaction(); }); - return a.get(); - }); - Assert.assertTrue(String.format("%s should have been strictly smaller than %s", after, before), before > after); - Assert.assertEquals(0, after); + after = countDiffs.call(); + if (after == 0) + return; + Thread.sleep(1000); + } + Assert.fail("Should have GC'd all in (way under) 60 cycles. Remaining: " + after); }); } } -} - +} \ No newline at end of file