This commit is contained in:
Alan Wang 2026-07-17 17:54:56 -07:00
parent 85d29aff2f
commit 9a67548bd8
2 changed files with 26 additions and 10 deletions

View File

@ -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<Long, CoordinatedTransfer> coordinating = new ConcurrentHashMap<>();
public final Map<Long, CoordinatedTransfer> coordinating = new ConcurrentHashMap<>();
// Added when we have a streamed SSTable in our pending directory
private final Map<TimeUUID, PendingLocalTransfer> local = new ConcurrentHashMap<>();
public final Map<TimeUUID, PendingLocalTransfer> 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
{

View File

@ -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<SSTableReader> sstablesPriorToMove = new ArrayList<>(sstables.size());
Collection<SSTableReader> moved = new ArrayList<>(sstables.size());
Collection<Descriptor> 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);