diff --git a/CHANGES.txt b/CHANGES.txt index 22502533e9..eefd3077ad 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -90,6 +90,7 @@ Merged from 5.0: * Fix resource cleanup after SAI query timeouts (CASSANDRA-19177) * Suppress CVE-2023-6481 (CASSANDRA-19184) Merged from 4.1: + * 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 5348087279..248d82b2e4 100644 --- a/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosUncommittedIndex.java +++ b/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosUncommittedIndex.java @@ -64,6 +64,7 @@ import org.apache.cassandra.schema.Indexes; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.service.ClientState; import org.apache.cassandra.utils.CloseableIterator; +import org.apache.cassandra.utils.concurrent.OpOrder; import static java.util.Collections.singletonList; import static org.apache.cassandra.schema.SchemaConstants.SYSTEM_KEYSPACE_NAME; @@ -150,21 +151,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)