From e14816e244e595096ff1e7fdab829fee3031fe5e Mon Sep 17 00:00:00 2001 From: Benedict Elliott Smith Date: Tue, 16 Dec 2025 12:52:55 +0000 Subject: [PATCH] Journal reads must select segments before sstables to avoid compaction races Also Fix Cassandra: - In memory size calculation for CommandsForKey include Unmanaged - Accord load out-of-band cleanup should use SafeRedundantBefore ALso Improve Cassandra: - Report replay information on begin replay - Improve AccordService shutdown - Log command store RedundantBefore on shutdown - Segment compaction should wait for readOrder barrier to replace segments, for additional safety - Journal segments should share readOrder with sstables Also Improve Accord: - Iterate LocalListeners in order, so can query more effectively on node - Refine AbstractReplay.minReplay/shouldReplay patch by Benedict; reviewed by Alex Petrov for CASSANDRA-21804 --- modules/accord | 2 +- .../db/virtual/AccordDebugKeyspace.java | 17 ++- .../org/apache/cassandra/journal/Journal.java | 45 ++++---- .../apache/cassandra/journal/Segments.java | 2 +- .../cassandra/service/accord/AccordCache.java | 2 +- .../service/accord/AccordCommandStore.java | 42 +++++--- .../service/accord/AccordCommandStores.java | 39 ++++--- .../service/accord/AccordJournal.java | 15 ++- .../service/accord/AccordJournalTable.java | 100 ++++++++++-------- .../service/accord/AccordObjectSizes.java | 10 ++ .../accord/AccordSegmentCompactor.java | 2 + .../service/accord/AccordService.java | 47 ++------ .../service/accord/AccordTracing.java | 2 +- .../service/accord/CommandsForRanges.java | 2 +- .../accord/AccordJournalCompactionTest.java | 1 - .../test/AccordJournalSimulationTest.java | 4 +- .../apache/cassandra/journal/JournalTest.java | 5 +- 17 files changed, 194 insertions(+), 143 deletions(-) diff --git a/modules/accord b/modules/accord index f6b0a6998f..f09a12da76 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit f6b0a6998faca767e6951976097dec704c306b0e +Subproject commit f09a12da76bbc195ceb05ad859912aeb0a432dda diff --git a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java index c3f08de148..a670bdb72b 100644 --- a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java +++ b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java @@ -692,7 +692,7 @@ public class AccordDebugKeyspace extends VirtualKeyspace " waiting_until text,\n" + " waiter 'TxnIdUtf8Type',\n" + " PRIMARY KEY (command_store_id, waiting_on, waiting_until, waiter)" + - ')', Int32Type.instance), FAIL, UNSORTED, ASC); + ')', Int32Type.instance), FAIL, ASC, ASC); } @Override @@ -726,10 +726,15 @@ public class AccordDebugKeyspace extends VirtualKeyspace LocalListeners listeners = safeStore.commandStore().unsafeGetListeners(); for (LocalListeners.TxnListener listener : listeners.txnListeners()) { - rows.add(listener.waitingOn.toString(), listener.awaitingStatus.name(), listener.waiter.toString()) + rows.add(listener.waitingOn.toString(), ordered(listener.awaitingStatus), listener.waiter.toString()) .eagerCollect(ignore -> {}); } } + + private String ordered(SaveStatus saveStatus) + { + return (saveStatus.ordinal() <= 9 ? "0" : "") + saveStatus.ordinal() + '_' + saveStatus.name(); + } } public static final class ListenersLocalTable extends AbstractLazyVirtualTable @@ -1796,7 +1801,10 @@ public class AccordDebugKeyspace extends VirtualKeyspace case TRY_EXECUTE: run(txnId, commandStoreId, safeStore -> { SafeCommand safeCommand = safeStore.unsafeGet(txnId); - Commands.maybeExecute(safeStore, safeCommand, safeCommand.current(), true, true, NotifyWaitingOnPlus.adapter(ignore -> {}, true, true)); + Command command = safeCommand.current(); + if (command.saveStatus() == SaveStatus.Applying) + return Commands.applyChain(safeStore, (Command.Executed) command); + Commands.maybeExecute(safeStore, safeCommand, command, true, true, NotifyWaitingOnPlus.adapter(ignore -> {}, true, true)); return AsyncChains.success(null); }); break; @@ -1896,7 +1904,8 @@ public class AccordDebugKeyspace extends VirtualKeyspace AccordService.getBlocking(accord.node() .commandStores() .forId(commandStoreId) - .chain(PreLoadContext.contextFor(txnId, TXN_OPS), apply)); + .chain(PreLoadContext.contextFor(txnId, TXN_OPS), apply) + .flatMap(i -> i)); } private void cleanup(TxnId txnId, int commandStoreId, Cleanup cleanup) diff --git a/src/java/org/apache/cassandra/journal/Journal.java b/src/java/org/apache/cassandra/journal/Journal.java index ecfca32a52..a79873b2af 100644 --- a/src/java/org/apache/cassandra/journal/Journal.java +++ b/src/java/org/apache/cassandra/journal/Journal.java @@ -121,7 +121,7 @@ public class Journal implements Shutdownable private final FlusherCallbacks flusherCallbacks; - final OpOrder readOrder = new OpOrder(); + final OpOrder readOrder; private class FlusherCallbacks implements Flusher.Callbacks { @@ -177,7 +177,8 @@ public class Journal implements Shutdownable Params params, KeySupport keySupport, ValueSerializer valueSerializer, - SegmentCompactor segmentCompactor) + SegmentCompactor segmentCompactor, + OpOrder readOrder) { this.name = name; this.directory = directory; @@ -185,6 +186,7 @@ public class Journal implements Shutdownable this.keySupport = keySupport; this.valueSerializer = valueSerializer; + this.readOrder = readOrder; this.metrics = new Metrics<>(name); this.flusherCallbacks = new FlusherCallbacks(); @@ -357,15 +359,25 @@ public class Journal implements Shutdownable return null; } - public void readAll(K id, RecordConsumer consumer) + public static void readAll(K id, RecordConsumer consumer, OpOrder.Group readGroup, Segments segments) { EntrySerializer.EntryHolder holder = new EntrySerializer.EntryHolder<>(); - try (OpOrder.Group group = readOrder.start()) + for (Segment segment : segments.allSorted(false)) { - for (Segment segment : segments.get().allSorted(false)) - { - segment.readAll(id, holder, consumer); - } + segment.readAll(id, holder, consumer); + } + } + + public void readAll(K id, RecordConsumer consumer, OpOrder.Group readGroup) + { + readAll(id, consumer, readGroup, segments.get()); + } + + public void readAll(K id, RecordConsumer consumer) + { + try (OpOrder.Group readGroup = readOrder.start()) + { + readAll(id, consumer, readGroup); } } @@ -449,18 +461,15 @@ public class Journal implements Shutdownable * @return true if the record was found, false otherwise */ @SuppressWarnings("unused") - public boolean readLast(K id, RecordConsumer consumer) + public static boolean readLast(K id, RecordConsumer consumer, OpOrder.Group readOrder, Segments segments) { - try (OpOrder.Group group = readOrder.start()) + for (Segment segment : segments.allSorted(false)) { - for (Segment segment : segments.get().allSorted(false)) - { - if (!segment.index().mayContainId(id)) - continue; + if (!segment.index().mayContainId(id)) + continue; - if (segment.readLast(id, consumer)) - return true; - } + if (segment.readLast(id, consumer)) + return true; } return false; } @@ -728,7 +737,7 @@ public class Journal implements Shutdownable } } - Segments segments() + public Segments segments() { return segments.get(); } diff --git a/src/java/org/apache/cassandra/journal/Segments.java b/src/java/org/apache/cassandra/journal/Segments.java index 7245029fea..effc879513 100644 --- a/src/java/org/apache/cassandra/journal/Segments.java +++ b/src/java/org/apache/cassandra/journal/Segments.java @@ -31,7 +31,7 @@ import org.apache.cassandra.utils.concurrent.Refs; *

* TODO (performance, expected): an interval/range structure for StaticSegment lookup based on min/max key bounds */ -class Segments +public class Segments { private final Long2ObjectHashMap> segments; private SortedArrayList> sorted; diff --git a/src/java/org/apache/cassandra/service/accord/AccordCache.java b/src/java/org/apache/cassandra/service/accord/AccordCache.java index 6c6eb167b8..040baf11b2 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCache.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCache.java @@ -1321,7 +1321,7 @@ public class AccordCache implements CacheSize try (DataInputBuffer buf = new DataInputBuffer(buffer, false)) { builder.deserializeNext(buf, Version.LATEST); - return builder.construct(commandStore.unsafeGetRedundantBefore()); + return builder.construct(commandStore.safeGetRedundantBefore()); } catch (UnknownTableException e) { diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index 2e34080f99..3614bb65b2 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -317,7 +317,7 @@ public class AccordCommandStore extends CommandStore CommandsForKey cfk = CommandsForKeyAccessor.load(id, (TokenKey) key); if (cfk == null) return null; - RedundantBefore.QuickBounds bounds = unsafeGetRedundantBefore().get(key); + RedundantBefore.QuickBounds bounds = safeGetRedundantBefore().get(key); if (bounds == null) return cfk; // TODO (required): I don't think this should be possible? but we hit it on some test return cfk.withGcBeforeAtLeast(bounds.gcBefore, false); @@ -420,7 +420,7 @@ public class AccordCommandStore extends CommandStore @VisibleForTesting public Command loadCommand(TxnId txnId) { - return journal.loadCommand(id, txnId, unsafeGetRedundantBefore(), durableBefore()); + return journal.loadCommand(id, txnId, safeGetRedundantBefore(), durableBefore()); } @VisibleForTesting @@ -446,12 +446,12 @@ public class AccordCommandStore extends CommandStore public Command.Minimal loadMinimal(TxnId txnId) { - return journal.loadMinimal(id, txnId, unsafeGetRedundantBefore(), durableBefore()); + return journal.loadMinimal(id, txnId, safeGetRedundantBefore(), durableBefore()); } public Command.MinimalWithDeps loadMinimalWithDeps(TxnId txnId) { - return journal.loadMinimalWithDeps(id, txnId, unsafeGetRedundantBefore(), durableBefore()); + return journal.loadMinimalWithDeps(id, txnId, safeGetRedundantBefore(), durableBefore()); } public AccordCompactionInfo getCompactionInfo() @@ -465,6 +465,11 @@ public class AccordCommandStore extends CommandStore return new AccordCompactionInfo(id, redundantBefore, ranges, tableId); } + public final RedundantBefore safeGetRedundantBefore() + { + return safeRedundantBefore.redundantBefore; + } + public RangeSearcher rangeSearcher() { return rangeSearcher; @@ -472,10 +477,10 @@ public class AccordCommandStore extends CommandStore public AccordCommandStoreReplayer replayer() { - boolean replayOnlyDurable = true; + boolean replayOnlyNonDurable = true; if (journal instanceof AccordJournal) - replayOnlyDurable = ((AccordJournal)journal).configuration().replayMode() == ONLY_NON_DURABLE; - return new AccordCommandStoreReplayer(this, replayOnlyDurable); + replayOnlyNonDurable = ((AccordJournal)journal).configuration().replayMode() == ONLY_NON_DURABLE; + return new AccordCommandStoreReplayer(this, replayOnlyNonDurable); } static final AtomicLong nextDurabilityLoggingId = new AtomicLong(); @@ -540,15 +545,22 @@ public class AccordCommandStore extends CommandStore super.unsafeUpsertRedundantBefore(addRedundantBefore); } + @VisibleForTesting + public void unsafeUpdateRangesForEpoch() + { + super.unsafeUpdateRangesForEpoch(); + safeRedundantBefore = new SafeRedundantBefore(0, unsafeGetRedundantBefore()); + } + public static class AccordCommandStoreReplayer extends AbstractReplayer { - private final AccordCommandStore store; + private final AccordCommandStore commandStore; private final boolean onlyNonDurable; - private AccordCommandStoreReplayer(AccordCommandStore store, boolean onlyNonDurable) + private AccordCommandStoreReplayer(AccordCommandStore commandStore, boolean onlyNonDurable) { - super(store.unsafeGetRedundantBefore()); - this.store = store; + super(commandStore, null); + this.commandStore = commandStore; this.onlyNonDurable = onlyNonDurable; } @@ -558,7 +570,7 @@ public class AccordCommandStore extends CommandStore if (onlyNonDurable && !maybeShouldReplay(txnId)) return AsyncChains.success(null); - return store.chain(PreLoadContext.contextFor(txnId, "Replay"), safeStore -> { + return commandStore.chain(PreLoadContext.contextFor(txnId, "Replay"), safeStore -> { if (onlyNonDurable && !shouldReplay(txnId, safeStore.unsafeGet(txnId).current().participants())) return null; @@ -574,12 +586,16 @@ public class AccordCommandStore extends CommandStore void maybeLoadRedundantBefore(RedundantBefore redundantBefore) { + Invariants.require(safeRedundantBefore == null); if (redundantBefore != null) { loadRedundantBefore(redundantBefore); - Invariants.require(safeRedundantBefore == null); safeRedundantBefore = new SafeRedundantBefore(0, redundantBefore); } + else + { + safeRedundantBefore = new SafeRedundantBefore(0, this.unsafeGetRedundantBefore()); + } } void maybeLoadBootstrapBeganAt(NavigableMap bootstrapBeganAt) diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java index ddc66f36e0..c8a05795d4 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStores.java @@ -20,6 +20,7 @@ package org.apache.cassandra.service.accord; import java.util.Arrays; import java.util.List; import java.util.concurrent.TimeUnit; +import java.util.stream.Stream; import accord.api.Agent; import accord.api.DataStore; @@ -33,6 +34,7 @@ import accord.local.ShardDistributor; import accord.utils.RandomSource; import org.apache.cassandra.cache.CacheSize; import org.apache.cassandra.concurrent.ScheduledExecutors; +import org.apache.cassandra.concurrent.Shutdownable; import org.apache.cassandra.config.AccordSpec.QueueShardModel; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.service.accord.AccordExecutor.AccordExecutorFactory; @@ -43,8 +45,9 @@ import static org.apache.cassandra.config.DatabaseDescriptor.getAccordQueueSubmi import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITHOUT_LOCK; import static org.apache.cassandra.service.accord.AccordExecutor.Mode.RUN_WITH_LOCK; import static org.apache.cassandra.service.accord.AccordExecutor.constant; +import static org.apache.cassandra.utils.Clock.Global.nanoTime; -public class AccordCommandStores extends CommandStores implements CacheSize +public class AccordCommandStores extends CommandStores implements CacheSize, Shutdownable { private final AccordExecutor[] executors; private final int mask; @@ -206,25 +209,37 @@ public class AccordCommandStores extends CommandStores implements CacheSize executor.waitForQuiescence(); } + @Override + public boolean isTerminated() + { + return Stream.of(executors).allMatch(AccordExecutor::isTerminated); + } + @Override public synchronized void shutdown() { super.shutdown(); for (AccordExecutor executor : executors) - { executor.shutdown(); - } + } + + @Override + public Object shutdownNow() + { + shutdown(); + return null; + } + + @Override + public boolean awaitTermination(long timeout, TimeUnit units) throws InterruptedException + { + long deadline = nanoTime() + units.toNanos(timeout); for (AccordExecutor executor : executors) { - try - { - executor.awaitTermination(1, TimeUnit.MINUTES); - } - catch (InterruptedException e) - { - throw new RuntimeException(e); - } + long wait = Math.max(1, deadline - nanoTime()); + if (!executor.awaitTermination(wait, TimeUnit.NANOSECONDS)) + return false; } - //TODO (expected): shutdown isn't useful by itself, we need a way to "wait" as well. Should be AutoCloseable or offer awaitTermination as well (think Shutdownable interface) + return true; } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournal.java b/src/java/org/apache/cassandra/service/accord/AccordJournal.java index 885bbfa961..2127f336d8 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournal.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournal.java @@ -38,6 +38,7 @@ import com.google.common.annotations.VisibleForTesting; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import accord.impl.AbstractReplayer; import accord.impl.CommandChange; import accord.impl.CommandChange.Field; import accord.local.Cleanup; @@ -109,6 +110,8 @@ import static accord.impl.CommandChange.toIterableSetFields; import static accord.impl.CommandChange.unsetIterable; import static accord.impl.CommandChange.validateFlags; import static accord.local.Cleanup.Input.FULL; +import static accord.local.RedundantStatus.Property.LOCALLY_DURABLE_TO_COMMAND_STORE; +import static accord.local.RedundantStatus.Property.LOCALLY_DURABLE_TO_DATA_STORE; import static org.apache.cassandra.service.accord.AccordJournalValueSerializers.DurableBeforeAccumulator; import static org.apache.cassandra.service.accord.JournalKey.Type.COMMAND_DIFF; import static org.apache.cassandra.service.accord.journal.AccordTopologyUpdate.Accumulator; @@ -157,7 +160,8 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier throw new UnsupportedOperationException(); } }, - compactor(cfs, userVersion)); + compactor(cfs, userVersion), + cfs.readOrdering); this.journalTable = new AccordJournalTable<>(journal, JournalKey.SUPPORT, cfs, userVersion); this.params = params; } @@ -617,20 +621,23 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier class ReplayStream implements Closeable { final CommandStore commandStore; - final Replayer replayer; + final AbstractReplayer replayer; final CloseableIterator> iter; JournalKey prev; public ReplayStream(CommandStore commandStore) { this.commandStore = commandStore; - this.replayer = commandStore.replayer(); + this.replayer = (AbstractReplayer) commandStore.replayer(); // Keys in the index are sorted by command store id, so index iteration will be sequential - this.iter = journalTable.keyIterator(new JournalKey(TxnId.NONE, COMMAND_DIFF, commandStore.id()), new JournalKey(TxnId.MAX.withoutNonIdentityFlags(), COMMAND_DIFF, commandStore.id()), false); + this.iter = journalTable.keyIterator(new JournalKey(replayer.minReplay.withoutNonIdentityFlags(), COMMAND_DIFF, commandStore.id()), new JournalKey(TxnId.MAX.withoutNonIdentityFlags(), COMMAND_DIFF, commandStore.id()), false); } boolean replay() { + logger.info("Beginning replay of {} with min={}, {}", commandStore, replayer.minReplay, + replayer.redundantBefore.map(b -> b == null ? null : b.maxBoundBoth(LOCALLY_DURABLE_TO_DATA_STORE, LOCALLY_DURABLE_TO_COMMAND_STORE), TxnId[]::new)); + JournalKey key; long[] segments; while (true) diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java b/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java index 40649cc432..e557bdd993 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java @@ -76,6 +76,7 @@ import org.apache.cassandra.journal.Journal; import org.apache.cassandra.journal.KeySupport; import org.apache.cassandra.journal.RecordConsumer; import org.apache.cassandra.journal.Segment; +import org.apache.cassandra.journal.Segments; import org.apache.cassandra.schema.ColumnMetadata; import org.apache.cassandra.service.RetryStrategy; import org.apache.cassandra.service.accord.AccordKeyspace.JournalColumns; @@ -347,21 +348,26 @@ public class AccordJournalTable implements RangeSearche public void readAll(K key, RecordConsumer reader) { - try (TableKeyIterator table = readAllFromTable(key)) + try (OpOrder.Group readOrder = cfs.readOrdering.start()) { - boolean hasTableData = table.advance(); - long minSegment = hasTableData ? table.segment : Long.MIN_VALUE; - // First, read all journal entries newer than anything flushed into sstables - journal.readAll(key, (segment, position, key1, buffer, userVersion) -> { - if (segment > minSegment) - reader.accept(segment, position, key1, buffer, userVersion); - }); - - // Then, read SSTables - while (hasTableData) + // SELECT segments first, to avoid missing segments due to races compacting segment->sstable + Segments segments = journal.segments(); + try (TableKeyIterator table = readAllFromTable(key, readOrder)) { - reader.accept(table.segment, table.offset, key, table.value, table.userVersion); - hasTableData = table.advance(); + boolean hasTableData = table.advance(); + long minSegment = hasTableData ? table.segment : Long.MIN_VALUE; + // First, read all journal entries newer than anything flushed into sstables + Journal.readAll(key, (segment, position, key1, buffer, userVersion) -> { + if (segment > minSegment) + reader.accept(segment, position, key1, buffer, userVersion); + }, readOrder, segments); + + // Then, read SSTables + while (hasTableData) + { + reader.accept(table.segment, table.offset, key, table.value, table.userVersion); + hasTableData = table.advance(); + } } } } @@ -373,32 +379,36 @@ public class AccordJournalTable implements RangeSearche public void readLast(K key, RecordConsumer reader) { - try (TableKeyIterator table = readAllFromTable(key)) + try (OpOrder.Group readOrder = cfs.readOrdering.start()) { - boolean hasTableData = table.advance(); - long minSegment = hasTableData ? table.segment : Long.MIN_VALUE; - - class JournalReader implements RecordConsumer + Segments segments = journal.segments(); + try (TableKeyIterator table = readAllFromTable(key, readOrder)) { - boolean read; - @Override - public void accept(long segment, int position, K key, ByteBuffer buffer, int userVersion) + boolean hasTableData = table.advance(); + long minSegment = hasTableData ? table.segment : Long.MIN_VALUE; + + class JournalReader implements RecordConsumer { - if (segment > minSegment) + boolean read; + @Override + public void accept(long segment, int position, K key, ByteBuffer buffer, int userVersion) { - reader.accept(segment, position, key, buffer, userVersion); - read = true; + if (segment > minSegment) + { + reader.accept(segment, position, key, buffer, userVersion); + read = true; + } } } + + // First, read all journal entries newer than anything flushed into sstables + JournalReader journalReader = new JournalReader(); + Journal.readLast(key, journalReader, readOrder, segments); + + // Then, read SSTables, if we haven't found a record already + if (hasTableData && !journalReader.read) + reader.accept(table.segment, table.offset, key, table.value, table.userVersion); } - - // First, read all journal entries newer than anything flushed into sstables - JournalReader journalReader = new JournalReader(); - journal.readLast(key, journalReader); - - // Then, read SSTables, if we haven't found a record already - if (hasTableData && !journalReader.read) - reader.accept(table.segment, table.offset, key, table.value, table.userVersion); } } @@ -408,19 +418,17 @@ public class AccordJournalTable implements RangeSearche final K key; final List unmerged; final UnfilteredRowIterator merged; - final OpOrder.Group readOrder; long segment; int offset; ByteBuffer value; int userVersion; - TableKeyIterator(K key, List unmerged, UnfilteredRowIterator merged, OpOrder.Group readOrder) + TableKeyIterator(K key, List unmerged, UnfilteredRowIterator merged) { this.key = key; this.unmerged = unmerged; this.merged = merged; - this.readOrder = readOrder; } @Override @@ -455,16 +463,14 @@ public class AccordJournalTable implements RangeSearche @Override public void close() { - readOrder.close(); if (merged != null) merged.close(); } } - private TableKeyIterator readAllFromTable(K key) + private TableKeyIterator readAllFromTable(K key, OpOrder.Group readOrder) { DecoratedKey pk = JournalColumns.decorate(key); - OpOrder.Group readOrder = cfs.readOrdering.start(); List iters = new ArrayList<>(3); try { @@ -479,11 +485,10 @@ public class AccordJournalTable implements RangeSearche iters.add(iter); } - return new TableKeyIterator(key, iters, iters.isEmpty() ? null : UnfilteredRowIterators.merge(iters), readOrder); + return new TableKeyIterator(key, iters, iters.isEmpty() ? null : UnfilteredRowIterators.merge(iters)); } catch (Throwable t) { - readOrder.close(); for (UnfilteredRowIterator iter : iters) { try { iter.close(); } @@ -496,7 +501,10 @@ public class AccordJournalTable implements RangeSearche @SuppressWarnings("resource") // Auto-closeable iterator will release related resources public CloseableIterator> keyIterator(@Nullable K min, @Nullable K max, boolean includeActive) { - return new JournalAndTableKeyIterator(min, max, includeActive); + try (OpOrder.Group readOrder = cfs.readOrdering.start()) + { + return new JournalAndTableKeyIterator(min, max, includeActive); + } } private class TableIterator extends AbstractIterator implements CloseableIterator @@ -551,13 +559,19 @@ public class AccordJournalTable implements RangeSearche private class JournalAndTableKeyIterator extends AbstractIterator> implements CloseableIterator> { - final TableIterator tableIterator; final Journal.SegmentKeyIterator journalIterator; + final TableIterator tableIterator; private JournalAndTableKeyIterator(K min, K max, boolean includeActive) { - this.tableIterator = new TableIterator(min, max); + // We must initialise journal reader first, else we may race with segment->table compaction and miss some data + // that is, the following sequence could happen: + // - Select sstables to read + // - Segments compacted; segments removed and sstables added + // - Segment iterator created + // TODO (expected): segments should be sstables on creation this.journalIterator = journal.segmentKeyIterator(min, max, includeActive ? ignore -> true : Segment::isStatic); + this.tableIterator = new TableIterator(min, max); } K prevFromTable = null; diff --git a/src/java/org/apache/cassandra/service/accord/AccordObjectSizes.java b/src/java/org/apache/cassandra/service/accord/AccordObjectSizes.java index 31a2bf942f..d6955ec44c 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordObjectSizes.java +++ b/src/java/org/apache/cassandra/service/accord/AccordObjectSizes.java @@ -417,6 +417,7 @@ public class AccordObjectSizes private static long EMPTY_CFK_SIZE = measure(new CommandsForKey(null)); private static long EMPTY_INFO_SIZE = measure(CommandsForKey.NO_INFO); + private static long EMPTY_UNMANAGED_SIZE = measure(new CommandsForKey.Unmanaged(null, TxnId.NONE, TxnId.NONE)); private static long EMPTY_INFO_EXTRA_ADDITIONAL_SIZE = measure(TxnInfo.create(TxnId.NONE, ACCEPTED, false, TxnId.NONE, NO_TXNIDS, Ballot.MAX)) - EMPTY_INFO_SIZE; public static long commandsForKey(CommandsForKey cfk) { @@ -438,6 +439,15 @@ public class AccordObjectSizes size += ballot(infoExtra.ballot); } } + size += ObjectSizes.sizeOfReferenceArray(cfk.unmanagedCount()); + size += cfk.unmanagedCount() * EMPTY_UNMANAGED_SIZE; + size += cfk.unmanagedCount() * TIMESTAMP_SIZE; + for (int i = 0 ; i < cfk.unmanagedCount() ; ++i) + { + CommandsForKey.Unmanaged unmanaged = cfk.getUnmanaged(i); + if (unmanaged.waitingUntil != unmanaged.txnId) + size += TIMESTAMP_SIZE; + } return size; } } diff --git a/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java b/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java index 5f6c841a49..7cd047bc0e 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java @@ -57,6 +57,8 @@ public class AccordSegmentCompactor extends AbstractAccordSegmentCompactor cfs.addSSTables(writer.finish(true)); writer.close(); writer = null; + // await reads to complete before swapping segments, to provide additional guarantees against reading incomplete data + cfs.readOrdering.awaitNewBarrier(); } @Override diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index e9874fae0a..124eb3761c 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -134,6 +134,8 @@ import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; import static accord.api.Journal.TopologyUpdate; import static accord.api.ProtocolModifiers.Toggles.FastExec.MAY_BYPASS_SAFESTORE; import static accord.impl.progresslog.DefaultProgressLog.ModeFlag.CATCH_UP; +import static accord.local.RedundantStatus.Property.LOCALLY_DURABLE_TO_COMMAND_STORE; +import static accord.local.RedundantStatus.Property.LOCALLY_DURABLE_TO_DATA_STORE; import static accord.local.durability.DurabilityService.SyncLocal.Self; import static accord.local.durability.DurabilityService.SyncRemote.All; import static accord.messages.SimpleReply.Ok; @@ -252,7 +254,6 @@ public class AccordService implements IAccordService, Shutdownable private enum State { INIT, STARTED, SHUTTING_DOWN, SHUTDOWN } private final Node node; - private final Shutdownable nodeShutdown; private final AccordMessageSink messageSink; private final AccordEndpointMapper endpointMapper; private final AccordTopologyService topologyService; @@ -435,7 +436,6 @@ public class AccordService implements IAccordService, Shutdownable new AccordInteropFactory(endpointMapper), journal.durableBeforePersister(), journal); - this.nodeShutdown = toShutdownable(node); this.requestHandler = new AccordVerbHandler<>(node, endpointMapper); this.responseHandler = new AccordResponseVerbHandler<>(callbacks, endpointMapper); } @@ -1042,7 +1042,7 @@ public class AccordService implements IAccordService, Shutdownable private List shutdownableSubsystems() { - return Arrays.asList(scheduler, nodeShutdown, journal, topologyService); + return Arrays.asList((AccordCommandStores)node.commandStores(), journal, topologyService, scheduler); } @VisibleForTesting @@ -1051,6 +1051,10 @@ public class AccordService implements IAccordService, Shutdownable { if (!ExecutorUtils.shutdownThenWait(shutdownableSubsystems(), timeout, unit)) logger.error("One or more subsystems did not shut down cleanly."); + + node.commandStores().forAllUnsafe(commandStore -> { + logger.info("{} stopping with durability: {}", commandStore, commandStore.unsafeGetRedundantBefore().map(b -> b == null ? null : b.maxBoundBoth(LOCALLY_DURABLE_TO_DATA_STORE, LOCALLY_DURABLE_TO_COMMAND_STORE), TxnId[]::new)); + }); } @Override @@ -1119,43 +1123,6 @@ public class AccordService implements IAccordService, Shutdownable sink.respond(Ok, message); } - private static Shutdownable toShutdownable(Node node) - { - return new Shutdownable() { - private volatile boolean isShutdown = false; - - @Override - public boolean isTerminated() - { - // we don't know about terminiated... so settle for shutdown! - return isShutdown; - } - - @Override - public void shutdown() - { - isShutdown = true; - node.shutdown(); - } - - @Override - public Object shutdownNow() - { - // node doesn't offer shutdownNow - shutdown(); - return null; - } - - @Override - public boolean awaitTermination(long timeout, TimeUnit units) - { - // TODO (required): expose awaitTermination in Node - // node doesn't offer - return true; - } - }; - } - @VisibleForTesting public AccordEndpointMapper endpointMapper() { diff --git a/src/java/org/apache/cassandra/service/accord/AccordTracing.java b/src/java/org/apache/cassandra/service/accord/AccordTracing.java index c1a69d568d..eb1b845c92 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordTracing.java +++ b/src/java/org/apache/cassandra/service/accord/AccordTracing.java @@ -761,7 +761,7 @@ public class AccordTracing extends AccordCoordinatorMetrics.Listener public Tracing trace(TxnId txnId, @Nullable Participants participants, CoordinationKind kind) { if (kind == CoordinationKind.FetchDurableBefore) - return (cs, msg) -> logger.info("Catchup/FetchDurableBefore: {}", msg); + return (cs, msg) -> logger.info("Catchup/FetchDurableBefore: {}", msg.length() <= 100 ? msg : msg.substring(0, 100)); if (!txnIdMap.containsKey(txnId) && null == maybeTrace(txnId, participants, kind, NewOrFailure.NEW, traceNewPatterns)) return null; diff --git a/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java b/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java index 8aa9f3dd88..f831e82175 100644 --- a/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java +++ b/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java @@ -353,7 +353,7 @@ public class CommandsForRanges extends TreeMap implements Co } } - public CommandsForRanges.Loader loader(@Nullable TxnId primaryTxnId, LoadKeysFor loadKeysFor, Unseekables keysOrRanges) + public CommandsForRanges.Loader loader(TxnId primaryTxnId, LoadKeysFor loadKeysFor, Unseekables keysOrRanges) { RedundantBefore redundantBefore = commandStore.unsafeGetRedundantBefore(); MaxDecidedRX maxDecidedRX = commandStore.unsafeGetMaxDecidedRX(); diff --git a/test/distributed/org/apache/cassandra/service/accord/AccordJournalCompactionTest.java b/test/distributed/org/apache/cassandra/service/accord/AccordJournalCompactionTest.java index 8471fdd829..2585a8c87b 100644 --- a/test/distributed/org/apache/cassandra/service/accord/AccordJournalCompactionTest.java +++ b/test/distributed/org/apache/cassandra/service/accord/AccordJournalCompactionTest.java @@ -125,7 +125,6 @@ public class AccordJournalCompactionTest DurableBefore addDurableBefore = durableBeforeGen.next(rs); // TODO: improve redundant before generator and re-enable // updates.addRedundantBefore = redundantBeforeGen.next(rs); -// updates.newRedundantBefore = redundantBefore = RedundantBefore.merge(redundantBefore, updates.addRedundantBefore); updates.newSafeToRead = safeToReadGen.next(rs); updates.newRangesForEpoch = rangesForEpochGen.next(rs); diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java b/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java index 36040e115c..d220374cab 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java @@ -51,6 +51,7 @@ import org.apache.cassandra.journal.SegmentCompactor; import org.apache.cassandra.journal.ValueSerializer; import org.apache.cassandra.utils.Isolated; import org.apache.cassandra.utils.concurrent.CountDownLatch; +import org.apache.cassandra.utils.concurrent.OpOrder; public class AccordJournalSimulationTest extends SimulationTestBase { @@ -80,7 +81,8 @@ public class AccordJournalSimulationTest extends SimulationTestBase spec, new IdentityKeySerializer(), new IdentityValueSerializer(), - SegmentCompactor.noop()); + SegmentCompactor.noop(), + new OpOrder()); }), () -> check()); } diff --git a/test/unit/org/apache/cassandra/journal/JournalTest.java b/test/unit/org/apache/cassandra/journal/JournalTest.java index de6848b289..62342777a7 100644 --- a/test/unit/org/apache/cassandra/journal/JournalTest.java +++ b/test/unit/org/apache/cassandra/journal/JournalTest.java @@ -29,6 +29,7 @@ import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.io.util.File; import org.apache.cassandra.utils.TimeUUID; +import org.apache.cassandra.utils.concurrent.OpOrder; import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID; import static org.junit.Assert.assertEquals; @@ -49,7 +50,7 @@ public class JournalTest directory.deleteRecursiveOnExit(); Journal journal = - new Journal<>("TestJournal", directory, TestParams.INSTANCE, TimeUUIDKeySupport.INSTANCE, LongSerializer.INSTANCE, SegmentCompactor.noop()); + new Journal<>("TestJournal", directory, TestParams.INSTANCE, TimeUUIDKeySupport.INSTANCE, LongSerializer.INSTANCE, SegmentCompactor.noop(), new OpOrder()); journal.start(); @@ -70,7 +71,7 @@ public class JournalTest journal.shutdown(); - journal = new Journal<>("TestJournal", directory, TestParams.INSTANCE, TimeUUIDKeySupport.INSTANCE, LongSerializer.INSTANCE, SegmentCompactor.noop()); + journal = new Journal<>("TestJournal", directory, TestParams.INSTANCE, TimeUUIDKeySupport.INSTANCE, LongSerializer.INSTANCE, SegmentCompactor.noop(), new OpOrder()); journal.start(); assertEquals(1L, (long) journal.readLast(id1));