From 9a67548bd852bbb4a509d90b673b1a79b24fffc2 Mon Sep 17 00:00:00 2001 From: Alan Wang Date: Fri, 17 Jul 2026 17:54:56 -0700 Subject: [PATCH] WIP --- .../service/accord/LocalTransfers.java | 14 +++++------- .../service/accord/PendingLocalTransfer.java | 22 ++++++++++++++++++- 2 files changed, 26 insertions(+), 10 deletions(-) diff --git a/src/java/org/apache/cassandra/service/accord/LocalTransfers.java b/src/java/org/apache/cassandra/service/accord/LocalTransfers.java index 18f7601897..9ff16fd3f6 100644 --- a/src/java/org/apache/cassandra/service/accord/LocalTransfers.java +++ b/src/java/org/apache/cassandra/service/accord/LocalTransfers.java @@ -30,8 +30,6 @@ import com.google.common.base.Preconditions; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import accord.utils.Invariants; - import org.apache.cassandra.concurrent.ExecutorPlus; import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.io.util.File; @@ -52,10 +50,10 @@ public class LocalTransfers private final ReadWriteLock lock = new ReentrantReadWriteLock(); // SSTable imports that we are coordinating - private final Map coordinating = new ConcurrentHashMap<>(); + public final Map coordinating = new ConcurrentHashMap<>(); // Added when we have a streamed SSTable in our pending directory - private final Map local = new ConcurrentHashMap<>(); + public final Map local = new ConcurrentHashMap<>(); final ExecutorPlus executor = executorFactory().pooled("LocalTrackedTransfers", Integer.MAX_VALUE); @@ -160,6 +158,7 @@ public class LocalTransfers logger.debug("Deleting pending transfer directory: {}", pendingDir); pendingDir.deleteRecursive(); } + local.remove(transfer.planId); } } @@ -218,23 +217,20 @@ public class LocalTransfers }); } + // This method will be called by every Accord command store executor of + // the ranges it intersects public void activatePendingTransfers(TxnRead.ImportMetadata metadata) { lock.readLock().lock(); try { - int activatedTransfer = 0; for (TimeUUID planId : metadata.getPlanIds()) { PendingLocalTransfer pendingLocalTransfer = local.get(planId); if (pendingLocalTransfer != null) - { - activatedTransfer += 1; pendingLocalTransfer.activate(); - } } - Invariants.require(activatedTransfer == 1); } finally { diff --git a/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java b/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java index d402c5a2d8..3cb35a32cc 100644 --- a/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java +++ b/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java @@ -29,6 +29,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.io.sstable.Descriptor; import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.io.util.File; import org.apache.cassandra.schema.TableId; @@ -79,9 +80,28 @@ public class PendingLocalTransfer File dst = cfs.getDirectories().getDirectoryForNewSSTables(); dst.createFileIfNotExists(); + Collection sstablesPriorToMove = new ArrayList<>(sstables.size()); Collection moved = new ArrayList<>(sstables.size()); + Collection movedDescriptors = new ArrayList<>(sstables.size()); for (SSTableReader sstable : sstables) - moved.add(SSTableReader.moveAndOpenSSTable(cfs, sstable.descriptor, cfs.getUniqueDescriptorFor(sstable.descriptor, dst), sstable.getComponents(), true)); + { + try + { + sstablesPriorToMove.add(sstable); + Descriptor newDescriptor = cfs.getUniqueDescriptorFor(sstable.descriptor, dst); + movedDescriptors.add(newDescriptor); + SSTableReader movedSSTable = SSTableReader.moveAndOpenSSTable(cfs, sstable.descriptor, newDescriptor, sstable.getComponents(), true); + moved.add(movedSSTable); + } + catch (Throwable t) + { + sstablesPriorToMove.forEach(s -> s.selfRef().release()); + logger.error("Failed importing sstables"); + for (Descriptor descriptor : movedDescriptors) + descriptor.getFormat().delete(descriptor); + throw new RuntimeException("Failed importing SSTables", t); + } + } // Add all SSTables atomically cfs.getTracker().addSSTables(moved);