This is an automated email from the ASF dual-hosted git repository.

bdeggleston pushed a commit to branch cep-45-mutation-tracking
in repository https://gitbox.apache.org/repos/asf/cassandra.git


The following commit(s) were added to refs/heads/cep-45-mutation-tracking by 
this push:
     new 8fda60880b CEP 45 - Fix tracked SSTable import/transfer bugs
8fda60880b is described below

commit 8fda60880b4f1e1329c8cb0999c049dd2e9cea43
Author: Alan Wang <[email protected]>
AuthorDate: Tue May 12 09:28:43 2026 -0700

    CEP 45 - Fix tracked SSTable import/transfer bugs
    
    - Filter interval-tree false positives: skip SSTables that intersect a
      shard range but contain no partitions within it, precomputing and
      reusing partition positions per SSTable
    - Register PendingLocalTransfer once in finished() instead of per-SSTable
      in received(), fixing an assertion failure when streaming >1 SSTable
    - Guard against null currentWriter in RangeAwareSSTableWriter.setOpenResult
    - Schedule local cleanup before notifyFailure so cleanup always runs on
      the failure path, and remove completed transfers from the local map
    - Adds dtests for the interval-tree false positive and multi-SSTable import
      cases.
    
    Patch by Alan Wang; Reviewed by Blake Eggleston for CASSANDRA-21470
---
 .../db/streaming/CassandraStreamReceiver.java      | 10 +--
 .../io/sstable/RangeAwareSSTableWriter.java        |  3 +-
 .../replication/TrackedImportTransfer.java         |  9 ++-
 .../replication/TrackedImportTransfers.java        | 14 +++-
 .../replication/TransferTrackingService.java       |  1 +
 .../test/tracking/TrackedImportFailureTest.java    |  7 +-
 .../test/tracking/TrackedImportTransferTest.java   | 92 ++++++++++++++++++++++
 .../tracking/TrackedRepairTransferTestBase.java    |  3 +
 .../test/tracking/TrackedTransferTestBase.java     | 17 ++--
 9 files changed, 136 insertions(+), 20 deletions(-)

diff --git 
a/src/java/org/apache/cassandra/db/streaming/CassandraStreamReceiver.java 
b/src/java/org/apache/cassandra/db/streaming/CassandraStreamReceiver.java
index 70f494b338..b1571db404 100644
--- a/src/java/org/apache/cassandra/db/streaming/CassandraStreamReceiver.java
+++ b/src/java/org/apache/cassandra/db/streaming/CassandraStreamReceiver.java
@@ -149,12 +149,6 @@ public class CassandraStreamReceiver implements 
StreamReceiver
         txn.update(finished);
         sstables.addAll(finished);
         receivedEntireSSTable = file.isEntireSSTable();
-
-        if (useTrackedTransferPath())
-        {
-            PendingLocalTransfer transfer = new 
PendingLocalTransfer(cfs.metadata().id, session.planId(), sstables);
-            MutationTrackingService.instance().received(transfer);
-        }
     }
 
     @Override
@@ -282,7 +276,11 @@ public class CassandraStreamReceiver implements 
StreamReceiver
 
                 // SSTables involved in a coordinated transfer become live 
when the transfer is activated
                 if (useTrackedTransferPath())
+                {
+                    PendingLocalTransfer transfer = new 
PendingLocalTransfer(cfs.metadata().id, session.planId(), sstables);
+                    MutationTrackingService.instance().received(transfer);
                     return;
+                }
 
                 if (session.streamOperation() == StreamOperation.BOOTSTRAP)
                 {
diff --git 
a/src/java/org/apache/cassandra/io/sstable/RangeAwareSSTableWriter.java 
b/src/java/org/apache/cassandra/io/sstable/RangeAwareSSTableWriter.java
index a4c15ff53e..97b6273443 100644
--- a/src/java/org/apache/cassandra/io/sstable/RangeAwareSSTableWriter.java
+++ b/src/java/org/apache/cassandra/io/sstable/RangeAwareSSTableWriter.java
@@ -146,7 +146,8 @@ public class RangeAwareSSTableWriter implements 
SSTableMultiWriter
     public SSTableMultiWriter setOpenResult(boolean openResult)
     {
         finishedWriters.forEach((w) -> w.setOpenResult(openResult));
-        currentWriter.setOpenResult(openResult);
+        if (currentWriter != null)
+            currentWriter.setOpenResult(openResult);
         return this;
     }
 
diff --git 
a/src/java/org/apache/cassandra/replication/TrackedImportTransfer.java 
b/src/java/org/apache/cassandra/replication/TrackedImportTransfer.java
index 6710317317..2b6d3eaeb0 100644
--- a/src/java/org/apache/cassandra/replication/TrackedImportTransfer.java
+++ b/src/java/org/apache/cassandra/replication/TrackedImportTransfer.java
@@ -88,19 +88,22 @@ public class TrackedImportTransfer extends 
CoordinatedTransfer
 
     final Collection<SSTableReader> sstables;
     private final ConsistencyLevel cl;
+    final Map<SSTableReader, List<SSTableReader.PartitionPositionBounds>> 
positionForSSTables;
 
     @VisibleForTesting
     TrackedImportTransfer(Range<Token> range, MutationId id)
     {
         super(id, null, range);
         this.sstables = Collections.emptyList();
+        this.positionForSSTables = Collections.emptyMap();
         this.cl = null;
     }
 
-    TrackedImportTransfer(String keyspace, Range<Token> range, Participants 
participants, Collection<SSTableReader> sstables, ConsistencyLevel cl, 
Supplier<MutationId> nextId)
+    TrackedImportTransfer(String keyspace, Range<Token> range, Participants 
participants, Collection<SSTableReader> sstables, Map<SSTableReader, 
List<SSTableReader.PartitionPositionBounds>> positionForSSTables, 
ConsistencyLevel cl, Supplier<MutationId> nextId)
     {
         super(nextId.get(), participants, keyspace, range);
         this.sstables = sstables;
+        this.positionForSSTables = positionForSSTables;
         this.cl = cl;
 
         ClusterMetadata cm = ClusterMetadata.current();
@@ -228,8 +231,8 @@ public class TrackedImportTransfer extends 
CoordinatedTransfer
             if (!purgeable)
                 return;
 
-            notifyFailure();
             TransferTrackingService.instance().scheduleCleanup();
+            notifyFailure();
         }
         catch (Throwable t)
         {
@@ -361,7 +364,7 @@ public class TrackedImportTransfer extends 
CoordinatedTransfer
         for (SSTableReader sstable : sstables)
         {
             List<Range<Token>> ranges = Collections.singletonList(range);
-            List<SSTableReader.PartitionPositionBounds> positions = 
sstable.getPositionsForRanges(ranges);
+            List<SSTableReader.PartitionPositionBounds> positions = 
positionForSSTables.get(sstable);
             long estimatedKeys = sstable.estimatedKeysForRanges(ranges);
             OutgoingStream stream = new 
CassandraOutgoingFile(StreamOperation.IMPORT, sstable.ref(), positions, ranges, 
estimatedKeys);
             plan.transferStreams(to, Collections.singleton(stream));
diff --git 
a/src/java/org/apache/cassandra/replication/TrackedImportTransfers.java 
b/src/java/org/apache/cassandra/replication/TrackedImportTransfers.java
index 250c240e0a..b2974df4c6 100644
--- a/src/java/org/apache/cassandra/replication/TrackedImportTransfers.java
+++ b/src/java/org/apache/cassandra/replication/TrackedImportTransfers.java
@@ -22,8 +22,11 @@ import java.io.IOException;
 import java.io.UncheckedIOException;
 import java.util.ArrayList;
 import java.util.Collection;
+import java.util.Collections;
+import java.util.HashMap;
 import java.util.Iterator;
 import java.util.List;
+import java.util.Map;
 
 import org.apache.cassandra.db.ConsistencyLevel;
 import org.apache.cassandra.db.lifecycle.SSTableIntervalTree;
@@ -67,10 +70,19 @@ public class TrackedImportTransfers implements 
Iterable<TrackedImportTransfer>
         shards.forEachShard(shard -> {
             Range<Token> range = shard.tokenRange();
             Collection<SSTableReader> sstablesForRange = 
intervals.search(Interval.create(range.left.minKeyBound(), 
range.right.maxKeyBound()));
+            List<Range<Token>> ranges = Collections.singletonList(range);
+            Map<SSTableReader, List<SSTableReader.PartitionPositionBounds>> 
positionForSSTables = new HashMap<>();
+            sstablesForRange.removeIf(sstable -> {
+                List<SSTableReader.PartitionPositionBounds> position = 
sstable.getPositionsForRanges(ranges);
+                if (!position.isEmpty())
+                    positionForSSTables.put(sstable, position);
+                return position.isEmpty();
+            });
+
             if (sstablesForRange.isEmpty())
                 return;
 
-            TrackedImportTransfer transfer = new 
TrackedImportTransfer(keyspace, range, shard.participants, sstablesForRange, 
cl, shard::nextId);
+            TrackedImportTransfer transfer = new 
TrackedImportTransfer(keyspace, range, shard.participants, sstablesForRange, 
positionForSSTables, cl, shard::nextId);
             transfers.add(transfer);
         });
         return new TrackedImportTransfers(transfers);
diff --git 
a/src/java/org/apache/cassandra/replication/TransferTrackingService.java 
b/src/java/org/apache/cassandra/replication/TransferTrackingService.java
index 5729a047f2..5f6c22e020 100644
--- a/src/java/org/apache/cassandra/replication/TransferTrackingService.java
+++ b/src/java/org/apache/cassandra/replication/TransferTrackingService.java
@@ -346,6 +346,7 @@ public class TransferTrackingService
                     pendingDir.deleteRecursive();
                 }
             }
+            local.remove(transfer.planId);
         }
         finally
         {
diff --git 
a/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedImportFailureTest.java
 
b/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedImportFailureTest.java
index 72b8eaabf9..ee897a3188 100644
--- 
a/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedImportFailureTest.java
+++ 
b/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedImportFailureTest.java
@@ -331,7 +331,9 @@ public class TrackedImportFailureTest extends 
TrackedTransferTestBase
             // Await cleanup of failed stream
             Uninterruptibles.sleepUninterruptibly(5, TimeUnit.SECONDS);
 
-            assertPendingDirs(cluster, (File pendingUuidDir) -> {
+            // We exclude the missed instance because the SSTables streamed to 
the pending directory
+            // are not linked to TransferTrackingService and hence cleanup 
does not know which SSTables to clean up
+            assertPendingDirs(cluster.stream().filter(instance -> instance != 
missed).collect(Collectors.toList()), (File pendingUuidDir) -> {
                 
Assertions.assertThat(pendingUuidDir.listUnchecked(File::isFile)).isEmpty();
             });
 
@@ -350,10 +352,11 @@ public class TrackedImportFailureTest extends 
TrackedTransferTestBase
             createSchema(cluster, keyspace);
 
             Iterable<IInvokableInstance> down = 
Collections.singleton(cluster.get(3));
-            Iterable<IInvokableInstance> up = cluster.stream().filter(instance 
-> instance != down).collect(Collectors.toList());
             for (IInvokableInstance instance : down)
                 instance.shutdown().get();
 
+            Iterable<IInvokableInstance> up = cluster.stream().filter(instance 
-> !instance.isShutdown()).collect(Collectors.toList());
+
             doImport(cluster, keyspace);
 
             cluster.get(3).startup();
diff --git 
a/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedImportTransferTest.java
 
b/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedImportTransferTest.java
index 3c279afcd5..9095d82db3 100644
--- 
a/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedImportTransferTest.java
+++ 
b/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedImportTransferTest.java
@@ -19,13 +19,18 @@
 package org.apache.cassandra.distributed.test.tracking;
 
 import java.io.IOException;
+import java.nio.file.Files;
 import java.time.Duration;
 import java.util.Collections;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
+import java.util.concurrent.TimeUnit;
 import java.util.function.Consumer;
 
+import com.google.common.util.concurrent.Uninterruptibles;
+
+import org.assertj.core.api.Assertions;
 import org.junit.AfterClass;
 import org.junit.BeforeClass;
 import org.junit.Test;
@@ -36,7 +41,9 @@ import org.apache.cassandra.db.ColumnFamilyStore;
 import org.apache.cassandra.distributed.Cluster;
 import org.apache.cassandra.distributed.api.ConsistencyLevel;
 import org.apache.cassandra.distributed.api.IInvokableInstance;
+import org.apache.cassandra.distributed.shared.AssertUtils;
 import org.apache.cassandra.distributed.test.sai.SAIUtil;
+import org.apache.cassandra.io.sstable.CQLSSTableWriter;
 import org.apache.cassandra.io.util.File;
 import org.apache.cassandra.tcm.ClusterMetadata;
 import org.apache.cassandra.tcm.ownership.DataPlacement;
@@ -199,4 +206,89 @@ public class TrackedImportTransferTest extends 
TrackedTransferTestBase
             });
         }
     }
+
+    @Test
+    public void importIntervalTreeFalsePositive() throws IOException
+    {
+        // See CASSANDRA-21470
+        String keyspace = "interval_tree_false_positive";
+        cluster.schemaChange("CREATE KEYSPACE " + keyspace + " WITH 
replication = {'class': 'SimpleStrategy', 'replication_factor': 3} AND 
replication_type='tracked';");
+        cluster.schemaChange("CREATE TABLE " + tableWithKeyspace(keyspace) + " 
(k BLOB PRIMARY KEY, v INT)");
+
+        String file = 
Files.createTempDirectory(TrackedTransferTestBase.class.getSimpleName()).toString();
+
+        CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
+                                                           .forTable("CREATE 
TABLE " + tableWithKeyspace(keyspace) + " (k BLOB PRIMARY KEY, v INT)")
+                                                           .inDirectory(file)
+                                                           .using("INSERT INTO 
" + tableWithKeyspace(keyspace) + " (k, v) " + "VALUES (?, ?)");
+
+        // For shard (-3074457345618258603,3074457345618258601], this SSTable
+        // intersects it, but does not contain any values in between the shard.
+        try (CQLSSTableWriter writer = builder.build())
+        {
+            writer.addRow(KEY_100, 1); // -4074457345618258601L
+            writer.addRow(KEY_300, 1); // 3074457345618258602L
+        }
+
+        // empty
+        assertLocalSelect(cluster, keyspace, AssertUtils::assertRows);
+
+        List<String> failed = cluster.get(1).callOnInstance(() -> {
+            ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(keyspace, 
TABLE);
+            Set<String> paths = Set.of(file);
+            logger.info("Importing SSTables {}", paths);
+            return cfs.importNewSSTables(paths, true, true, true, true, true, 
true, true);
+        });
+
+        // Sleep for a while to make sure import completes
+        Uninterruptibles.sleepUninterruptibly(3, TimeUnit.SECONDS);
+
+        Assertions.assertThat(failed).isEmpty();
+        assertLocalSelect(cluster, keyspace, rows -> assertRows(rows, 
row(KEY_100, 1), row(KEY_300, 1)));
+    }
+
+    @Test
+    public void importMoreThanOneSSTable() throws IOException
+    {
+        String keyspace = "import_more_than_one_sstable";
+        cluster.schemaChange("CREATE KEYSPACE " + keyspace + " WITH 
replication = {'class': 'SimpleStrategy', 'replication_factor': 3} AND 
replication_type='tracked';");
+        cluster.schemaChange("CREATE TABLE " + tableWithKeyspace(keyspace) + " 
(k int PRIMARY KEY, v INT)");
+
+        String file = 
Files.createTempDirectory(TrackedTransferTestBase.class.getSimpleName()).toString();
+
+        CQLSSTableWriter.Builder builder1 = CQLSSTableWriter.builder()
+                                                           .forTable("CREATE 
TABLE " + tableWithKeyspace(keyspace) + " (k int PRIMARY KEY, v INT)")
+                                                           .inDirectory(file)
+                                                           .using("INSERT INTO 
" + tableWithKeyspace(keyspace) + " (k, v) " + "VALUES (?, ?)");
+
+        try (CQLSSTableWriter writer = builder1.build())
+        {
+            writer.addRow(1, 1);
+        }
+
+        CQLSSTableWriter.Builder builder2 = CQLSSTableWriter.builder()
+                                                           .forTable("CREATE 
TABLE " + tableWithKeyspace(keyspace) + " (k int PRIMARY KEY, v INT)")
+                                                           .inDirectory(file)
+                                                           .using("INSERT INTO 
" + tableWithKeyspace(keyspace) + " (k, v) " + "VALUES (?, ?)");
+
+        try (CQLSSTableWriter writer = builder2.build())
+        {
+            writer.addRow(8, 1);
+        }
+
+        assertLocalSelect(cluster, keyspace, AssertUtils::assertRows);
+
+        List<String> failed = cluster.get(1).callOnInstance(() -> {
+            ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(keyspace, 
TABLE);
+            Set<String> paths = Set.of(file);
+            logger.info("Importing SSTables {}", paths);
+            return cfs.importNewSSTables(paths, true, true, true, true, true, 
true, true);
+        });
+
+        // Sleep for a while to make sure import completes
+        Uninterruptibles.sleepUninterruptibly(3, TimeUnit.SECONDS);
+
+        Assertions.assertThat(failed).isEmpty();
+        assertLocalSelect(cluster, keyspace, rows -> assertRows(rows, row(1, 
1), row(8, 1)));
+    }
 }
diff --git 
a/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedRepairTransferTestBase.java
 
b/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedRepairTransferTestBase.java
index 759b57a1fb..2f6c9eddd0 100644
--- 
a/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedRepairTransferTestBase.java
+++ 
b/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedRepairTransferTestBase.java
@@ -42,6 +42,9 @@ public abstract class TrackedRepairTransferTestBase extends 
TrackedTransferTestB
 
         IInvokableInstance coordinator = cluster.get(1);
         coordinator.executeInternal("INSERT INTO " + 
tableWithKeyspace(keyspace) + " (pk, v) VALUES (?, 1)", key);
+        coordinator.flush(keyspace);
+        coordinator.executeInternal("INSERT INTO " + 
tableWithKeyspace(keyspace) + " (pk, v) VALUES (?, 1)", key);
+        coordinator.flush(keyspace);
 
         // Write should only be present on instance 1
         cluster.forEach(instance -> {
diff --git 
a/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedTransferTestBase.java
 
b/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedTransferTestBase.java
index 513dcd35eb..171255aced 100644
--- 
a/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedTransferTestBase.java
+++ 
b/test/distributed/org/apache/cassandra/distributed/test/tracking/TrackedTransferTestBase.java
@@ -126,6 +126,10 @@ public abstract class TrackedTransferTestBase extends 
TestBaseImpl
     protected final static Token TOKEN_201 = new 
Murmur3Partitioner.LongToken(TOKEN_VALUE_201);
     protected final static ByteBuffer KEY_201 = 
Murmur3Partitioner.LongToken.keyForToken(TOKEN_201.getLongValue());
 
+    protected final static long TOKEN_VALUE_300 = 3074457345618258602L;
+    protected final static Token TOKEN_300 = new 
Murmur3Partitioner.LongToken(TOKEN_VALUE_300);
+    protected final static ByteBuffer KEY_300 = 
Murmur3Partitioner.LongToken.keyForToken(TOKEN_300.getLongValue());
+
     protected final static Range<Token> SHARD_ALIGNED_RANGE_2 = new 
Range<>(new Murmur3Partitioner.LongToken(TOKEN_VALUE_200 - 10), new 
Murmur3Partitioner.LongToken(TOKEN_VALUE_200 + 10));
 
     static
@@ -138,6 +142,9 @@ public abstract class TrackedTransferTestBase extends 
TestBaseImpl
 
         reversed = Murmur3Partitioner.instance.decorateKey(KEY_201);
         Assertions.assertThat(reversed.getToken()).isEqualTo(TOKEN_201);
+
+        reversed = Murmur3Partitioner.instance.decorateKey(KEY_300);
+        Assertions.assertThat(reversed.getToken()).isEqualTo(TOKEN_300);
     }
 
     protected static Cluster cluster() throws IOException
@@ -191,12 +198,12 @@ public abstract class TrackedTransferTestBase extends 
TestBaseImpl
         assertPendingDirs(validate, KEYSPACE, forPendingUuidDir);
     }
 
-    protected static void assertPendingDirs(Iterable<IInvokableInstance> 
validate, String keysapce, IIsolatedExecutor.SerializableConsumer<File> 
forPendingUuidDir)
+    protected static void assertPendingDirs(Iterable<IInvokableInstance> 
validate, String keyspace, IIsolatedExecutor.SerializableConsumer<File> 
forPendingUuidDir)
     {
         for (IInvokableInstance instance : validate)
         {
             instance.runOnInstance(() -> {
-                Set<File> allPendingDirs = 
ColumnFamilyStore.getIfExists(keysapce, 
TABLE).getDirectories().getPendingLocations();
+                Set<File> allPendingDirs = 
ColumnFamilyStore.getIfExists(keyspace, 
TABLE).getDirectories().getPendingLocations();
                 for (File pendingDir : allPendingDirs)
                 {
                     File[] pendingUuidDirs = 
pendingDir.listUnchecked(File::isDirectory);
@@ -315,7 +322,7 @@ public abstract class TrackedTransferTestBase extends 
TestBaseImpl
 
     protected static void doImport(Cluster cluster, IInvokableInstance target, 
Consumer<List<String>> onFailedDirs, String keyspace, @Nullable String 
createIndexCql) throws IOException
     {
-        String file = 
Files.createTempDirectory(MutationTrackingTest.class.getSimpleName()).toString();
+        String file = 
Files.createTempDirectory(TrackedTransferTestBase.class.getSimpleName()).toString();
 
         // Needs to run outside of instance executor because creates schema
         CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
@@ -357,10 +364,6 @@ public abstract class TrackedTransferTestBase extends 
TestBaseImpl
     {
         for (IInvokableInstance instance : validate)
         {
-            {
-                Object[][] rows = 
instance.executeInternal(withKeyspace("SELECT * FROM %s." + TABLE + " WHERE k = 
1", keyspace));
-                onRows.accept(rows);
-            }
             {
                 Object[][] rows = 
instance.executeInternal(withKeyspace("SELECT * FROM %s." + TABLE, keyspace));
                 onRows.accept(rows);


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to