From 2112c4c1c015db7e26edf4ba0c4eab6d0432fd56 Mon Sep 17 00:00:00 2001 From: Jon Haddad Date: Thu, 30 May 2024 08:24:22 -0700 Subject: [PATCH] Use OpOrder in repairIterator to ensure we don't lose memtables mid-paxos repair Patch by Jon Haddad; Reviewed by Blake Eggleston for CASSANDRA-19668 --- CHANGES.txt | 1 + .../uncommitted/PaxosUncommittedIndex.java | 30 +++++++++++-------- 2 files changed, 18 insertions(+), 13 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 5b41c58f71..847dce997a 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.1.6 + * Use OpOrder in repairIterator to ensure we don't lose memtables mid-paxos repair (Cassandra-19668) * Refresh stale paxos commit (CASSANDRA-19617) * Reduce info logging from automatic paxos repair (CASSANDRA-19445) * Support legacy plain_text_auth section in credentials file removed unintentionally (CASSANDRA-19498) diff --git a/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosUncommittedIndex.java b/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosUncommittedIndex.java index 5e3b5404e8..87a6ff9f81 100644 --- a/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosUncommittedIndex.java +++ b/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosUncommittedIndex.java @@ -49,6 +49,7 @@ import org.apache.cassandra.schema.IndexMetadata; import org.apache.cassandra.schema.Indexes; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.utils.CloseableIterator; +import org.apache.cassandra.utils.concurrent.OpOrder; import static java.util.Collections.*; import static org.apache.cassandra.schema.SchemaConstants.SYSTEM_KEYSPACE_NAME; @@ -135,21 +136,24 @@ public class PaxosUncommittedIndex implements Index, PaxosUncommittedTracker.Upd { Preconditions.checkNotNull(tableId); - View view = baseCfs.getTracker().getView(); - List memtables = view.flushingMemtables.isEmpty() - ? view.liveMemtables - : ImmutableList.builder().addAll(view.flushingMemtables).addAll(view.liveMemtables).build(); - - List dataRanges = ranges.stream().map(DataRange::forTokenRange).collect(Collectors.toList()); - List iters = new ArrayList<>(memtables.size() * ranges.size()); - - for (int j=0, jsize=dataRanges.size(); j memtables = view.flushingMemtables.isEmpty() + ? view.liveMemtables + : ImmutableList.builder().addAll(view.flushingMemtables).addAll(view.liveMemtables).build(); + + List dataRanges = ranges.stream().map(DataRange::forTokenRange).collect(Collectors.toList()); + List iters = new ArrayList<>(memtables.size() * ranges.size()); + + for (int j = 0, jsize = dataRanges.size(); j < jsize; j++) + { + for (int i = 0, isize = memtables.size(); i < isize; i++) + iters.add(memtables.get(i).partitionIterator(memtableColumnFilter, dataRanges.get(j), SSTableReadsListener.NOOP_LISTENER)); + } + return getPaxosUpdates(iters, tableId, false); + } } public CloseableIterator flushIterator(Memtable flushing)