From fd2e11f0d2a0b2f6254f9273109ed9e8e4aa7ab4 Mon Sep 17 00:00:00 2001 From: Benedict Elliott Smith Date: Thu, 10 Jul 2025 11:33:57 +0100 Subject: [PATCH] Accord Fixes: - WatermarkCollector should not report same closed/retired epoch N times - AccordSyncPropagator can merge pending requests and back-off retries - AccordCommandLoader should notify listeners - AccordSegmentCompactor should estimate number of keys to ensure bloom filters work - txn_blocked_by table should report what it can, not throw IllegalStateException - Permit uncompressed system tables patch by Benedict; reviewed by Alex Petrov for CASSANDRA-20754 --- modules/accord | 2 +- .../schema/CreateTableStatement.java | 2 +- .../db/virtual/AccordDebugKeyspace.java | 2 +- .../apache/cassandra/journal/OnDiskIndex.java | 5 ++ .../cassandra/journal/StaticSegment.java | 5 ++ .../AbstractAccordSegmentCompactor.java | 7 +- .../service/accord/AccordKeyspace.java | 5 +- .../accord/AccordSegmentCompactor.java | 4 +- .../service/accord/AccordSyncPropagator.java | 72 ++++++++++++++++++- .../service/accord/WatermarkCollector.java | 9 +-- .../accord/NemesisAccordSegmentCompactor.java | 2 +- 11 files changed, 97 insertions(+), 18 deletions(-) diff --git a/modules/accord b/modules/accord index 3c9b3077df..0d6157fc33 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 3c9b3077df1c74a14d8fe8dc58fb3f61b8257236 +Subproject commit 0d6157fc33dd16f1768030e66205e72fc1d7e9ac diff --git a/src/java/org/apache/cassandra/cql3/statements/schema/CreateTableStatement.java b/src/java/org/apache/cassandra/cql3/statements/schema/CreateTableStatement.java index 1f12366677..e452d36e11 100644 --- a/src/java/org/apache/cassandra/cql3/statements/schema/CreateTableStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/schema/CreateTableStatement.java @@ -173,7 +173,7 @@ public final class CreateTableStatement extends AlterSchemaStatement throw ire("read_repair must be set to 'NONE' for transiently replicated keyspaces"); } - if (!table.params.compression.isEnabled()) + if (!table.params.compression.isEnabled() && !SchemaConstants.isSystemKeyspace(table.keyspace)) Guardrails.uncompressedTablesEnabled.ensureEnabled(state); if (table.params.transactionalMode.accordIsEnabled && SchemaConstants.isSystemKeyspace(keyspaceName)) diff --git a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java index 4888c9d916..896c71c0ae 100644 --- a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java +++ b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java @@ -998,7 +998,7 @@ public class AccordDebugKeyspace extends VirtualKeyspace process(ds, commandStores, shard, processed, id, 0, id, Reason.Self, null); // everything was processed right? if (!shard.txns.isEmpty() && !shard.txns.keySet().containsAll(processed)) - throw new IllegalStateException("Skipped txns: " + Sets.difference(shard.txns.keySet(), processed)); + Invariants.expect(false, "Skipped txns: " + Sets.difference(shard.txns.keySet(), processed)); } return ds; diff --git a/src/java/org/apache/cassandra/journal/OnDiskIndex.java b/src/java/org/apache/cassandra/journal/OnDiskIndex.java index d82e5df5f3..5fae9255e3 100644 --- a/src/java/org/apache/cassandra/journal/OnDiskIndex.java +++ b/src/java/org/apache/cassandra/journal/OnDiskIndex.java @@ -199,6 +199,11 @@ final class OnDiskIndex extends Index return lastId; } + public int entryCount() + { + return entryCount; + } + @Override public long[] lookUp(K id) { diff --git a/src/java/org/apache/cassandra/journal/StaticSegment.java b/src/java/org/apache/cassandra/journal/StaticSegment.java index 281ed41178..7250659f85 100644 --- a/src/java/org/apache/cassandra/journal/StaticSegment.java +++ b/src/java/org/apache/cassandra/journal/StaticSegment.java @@ -239,6 +239,11 @@ public final class StaticSegment extends Segment return index; } + public int entryCount() + { + return index.entryCount(); + } + @Override boolean isActive() { diff --git a/src/java/org/apache/cassandra/service/accord/AbstractAccordSegmentCompactor.java b/src/java/org/apache/cassandra/service/accord/AbstractAccordSegmentCompactor.java index 954ae2a46c..b7b2539e05 100644 --- a/src/java/org/apache/cassandra/service/accord/AbstractAccordSegmentCompactor.java +++ b/src/java/org/apache/cassandra/service/accord/AbstractAccordSegmentCompactor.java @@ -84,7 +84,7 @@ public abstract class AbstractAccordSegmentCompactor implements SegmentCompac return false; } - abstract void initializeWriter(); + abstract void initializeWriter(int estimatedKeyCount); abstract SSTableTxnWriter writer(); abstract void finishAndAddWriter(); abstract Throwable cleanupWriter(Throwable t); @@ -99,9 +99,12 @@ public abstract class AbstractAccordSegmentCompactor implements SegmentCompac 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 + int estimatedKeyCount = 0; PriorityQueue> readers = new PriorityQueue<>(); for (StaticSegment segment : segments) { + estimatedKeyCount += segment.entryCount(); KeyOrderReader reader = segment.keyOrderReader(); if (reader.advance()) readers.add(reader); @@ -114,7 +117,7 @@ public abstract class AbstractAccordSegmentCompactor implements SegmentCompac if (readers.isEmpty()) return Collections.emptyList(); - initializeWriter(); + initializeWriter(estimatedKeyCount); JournalKey key = null; FlyweightImage builder = null; diff --git a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java index e597d7d415..81d97c4bb4 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java +++ b/src/java/org/apache/cassandra/service/accord/AccordKeyspace.java @@ -91,6 +91,7 @@ import org.apache.cassandra.io.sstable.SSTableReadsListener; import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.schema.ColumnMetadata; import org.apache.cassandra.schema.CompactionParams; +import org.apache.cassandra.schema.CompressionParams; import org.apache.cassandra.schema.IndexMetadata; import org.apache.cassandra.schema.Indexes; import org.apache.cassandra.schema.KeyspaceMetadata; @@ -144,8 +145,8 @@ public class AccordKeyspace + "user_version int," + "record blob," + "PRIMARY KEY((key), descriptor, offset)" - + ") WITH CLUSTERING ORDER BY (descriptor DESC, offset DESC)" + - " WITH compression = {'class':'NoopCompressor'};") + + ") WITH CLUSTERING ORDER BY (descriptor DESC, offset DESC);") + .compression(CompressionParams.NOOP) .compaction(CompactionParams.lcs(emptyMap())) .bloomFilterFpChance(0.01) .partitioner(new LocalPartitioner(BytesType.instance)); diff --git a/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java b/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java index 45633619a4..5f6c841a49 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSegmentCompactor.java @@ -37,12 +37,12 @@ public class AccordSegmentCompactor extends AbstractAccordSegmentCompactor } @Override - void initializeWriter() + void initializeWriter(int estimatedKeyCount) { Descriptor descriptor = cfs.newSSTableDescriptor(cfs.getDirectories().getDirectoryForNewSSTables()); SerializationHeader header = new SerializationHeader(true, cfs.metadata(), cfs.metadata().regularAndStaticColumns(), EncodingStats.NO_STATS); - this.writer = SSTableTxnWriter.create(cfs, descriptor, 0, 0, null, false, header); + this.writer = SSTableTxnWriter.create(cfs, descriptor, estimatedKeyCount, 0, null, false, header); } @Override diff --git a/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java b/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java index cf9e5db835..c8b5e75685 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java @@ -24,6 +24,7 @@ import java.util.Collection; import java.util.Collections; import java.util.HashSet; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; import com.google.common.collect.ImmutableSet; @@ -193,6 +194,7 @@ public class AccordSyncPropagator private final IFailureDetector failureDetector; private final ScheduledExecutorPlus scheduler; private final Listener listener; + private final ConcurrentHashMap retryingNotifications = new ConcurrentHashMap<>(); public AccordSyncPropagator(Node.Id localId, AccordEndpointMapper endpointMapper, MessageDelivery messagingService, IFailureDetector failureDetector, ScheduledExecutorPlus scheduler, @@ -304,6 +306,27 @@ public class AccordSyncPropagator }); } + private void scheduleRetry(Node.Id to, Notification notification) + { + Notification retry = new Notification(notification.epoch, notification.syncComplete, notification.closed, notification.retired, notification.attempts + 1); + RetryKey key = new RetryKey(to, notification.epoch); + retryingNotifications.compute(key, (k, cur) -> { + if (cur == null) + { + scheduler.schedule(() -> retry(k), Math.max(1, Math.min(15, retry.attempts)), TimeUnit.MINUTES); + return retry; + } + return cur.merge(retry); + }); + } + + private void retry(RetryKey key) + { + Notification retry = retryingNotifications.remove(key); + if (retry != null) + notify(key.to, retry); + } + private boolean notify(Node.Id to, Notification notification) { InetAddressAndPort toEp = endpointMapper.mappedEndpoint(to); @@ -335,7 +358,7 @@ public class AccordSyncPropagator @Override public void onFailure(InetAddressAndPort from, RequestFailure failure) { - scheduler.schedule(() -> AccordSyncPropagator.this.notify(to, notification), 1, TimeUnit.SECONDS); + scheduleRetry(to, notification); } @Override @@ -355,7 +378,7 @@ public class AccordSyncPropagator return true; } noSpamLogger.warn("Node{} is not alive, unable to notify of {}", to, notification); - scheduler.schedule(() -> notify(to, notification), 1, TimeUnit.MINUTES); + scheduleRetry(to, notification); return false; } messagingService.sendWithCallback(msg, toEp, cb); @@ -397,13 +420,30 @@ public class AccordSyncPropagator final long epoch; final Collection syncComplete; final Ranges closed, retired; + final int attempts; public Notification(long epoch, Collection syncComplete, Ranges closed, Ranges retired) + { + this(epoch, syncComplete, closed, retired, 0); + } + + public Notification(long epoch, Collection syncComplete, Ranges closed, Ranges retired, int attempts) { this.epoch = epoch; this.syncComplete = syncComplete; this.closed = closed; this.retired = retired; + this.attempts = attempts; + } + + Notification merge(Notification add) + { + Invariants.require(add.epoch == this.epoch); + Collection syncComplete = ImmutableSet.builder() + .addAll(this.syncComplete) + .addAll(add.syncComplete) + .build(); + return new Notification(epoch, syncComplete, closed.with(add.closed), retired.with(add.retired), Math.max(add.attempts, this.attempts)); } @Override @@ -417,4 +457,32 @@ public class AccordSyncPropagator '}'; } } + + static class RetryKey + { + final Node.Id to; + final long epoch; + + RetryKey(Node.Id id, long epoch) + { + to = id; + this.epoch = epoch; + } + + @Override + public int hashCode() + { + return to.id * 31 + (int)epoch; + } + + @Override + public boolean equals(Object obj) + { + if (!(obj instanceof RetryKey)) + return false; + + RetryKey that = (RetryKey) obj; + return that.epoch == this.epoch && that.to.equals(this.to); + } + } } diff --git a/src/java/org/apache/cassandra/service/accord/WatermarkCollector.java b/src/java/org/apache/cassandra/service/accord/WatermarkCollector.java index c4c3cb22cf..4ae51792ce 100644 --- a/src/java/org/apache/cassandra/service/accord/WatermarkCollector.java +++ b/src/java/org/apache/cassandra/service/accord/WatermarkCollector.java @@ -135,22 +135,19 @@ public class WatermarkCollector implements ConfigurationService.Listener MessageDelivery.RetryErrorMessage.EMPTY) .addCallback((m, fail) -> { if (fail != null) - { return; - } + Snapshot snapshot = m.payload; long minEpoch = configService.minEpoch(); for (Map.Entry e : snapshot.closed.entrySet()) { Ranges r = Ranges.of(e.getKey()); - for (long epoch = minEpoch; epoch <= e.getValue(); epoch++) - configService.receiveClosed(r, e.getValue()); + configService.receiveClosed(r, e.getValue()); } for (Map.Entry e : snapshot.retired.entrySet()) { Ranges r = Ranges.of(e.getKey()); - for (long epoch = minEpoch; epoch <= e.getValue(); epoch++) - configService.receiveRetired(r, e.getValue()); + configService.receiveRetired(r, e.getValue()); } for (Map.Entry e : snapshot.synced.entrySet()) { diff --git a/test/distributed/org/apache/cassandra/service/accord/NemesisAccordSegmentCompactor.java b/test/distributed/org/apache/cassandra/service/accord/NemesisAccordSegmentCompactor.java index b1e2407991..9c8a4f869c 100644 --- a/test/distributed/org/apache/cassandra/service/accord/NemesisAccordSegmentCompactor.java +++ b/test/distributed/org/apache/cassandra/service/accord/NemesisAccordSegmentCompactor.java @@ -63,7 +63,7 @@ public class NemesisAccordSegmentCompactor extends AbstractAccordSegmentCompa } @Override - void initializeWriter() + void initializeWriter(int estimatedKeyCount) { for (int i = 0; i < writers.length; i++) {