This commit is contained in:
Alan Wang 2026-07-20 13:59:06 -07:00
parent 9a67548bd8
commit 1b914f1a8e
4 changed files with 157 additions and 110 deletions

View File

@ -435,7 +435,7 @@ public class CoordinatedTransfer
public boolean equals(SingleTransferResult that)
{
return state == that.state && planId.equals(that.planId);
return state == that.state && (planId == null && that.planId == null || planId.equals(that.planId));
}
@Override

View File

@ -18,7 +18,9 @@
package org.apache.cassandra.service.accord;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
@ -31,7 +33,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.concurrent.ExecutorPlus;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.net.IVerbHandler;
import org.apache.cassandra.net.MessagingService;
@ -94,69 +95,35 @@ public class LocalTransfers
}
}
private void cleanupCoordinatedTransfer(CoordinatedTransfer transfer)
private void purge(TimeUUID timeUUID)
{
lock.writeLock().lock();
try
{
purge(transfer);
}
finally
{
lock.writeLock().unlock();
}
}
private void cleanupPendingLocalTransfer(TimeUUID timeUUID)
{
lock.writeLock().lock();
try
{
purge(local.get(timeUUID));
}
finally
{
lock.writeLock().unlock();
}
}
private void purge(TransferFailed failed)
{
lock.writeLock().lock();
try
{
PendingLocalTransfer pending = local.get(failed.planId);
if (pending == null)
PendingLocalTransfer transfer = local.get(timeUUID);
if (transfer == null)
{
logger.warn("Cannot purge unknown local pending transfer {}", failed);
logger.warn("Cannot purge unknown local pending transfer {}", transfer);
return;
}
purge(pending);
}
finally
{
lock.writeLock().unlock();
}
}
private void purge(PendingLocalTransfer transfer)
{
logger.info("Cleaning up pending transfer {}", transfer);
lock.writeLock().lock();
try
{
logger.info("Cleaning up pending transfer {}", transfer);
// Delete the entire pending transfer directory /pending/<planId>/
if (!transfer.sstables.isEmpty())
{
SSTableReader sstable = transfer.sstables.iterator().next();
File pendingDir = sstable.descriptor.directory;
Set<File> pendingDirs = new HashSet<>();
transfer.sstables.forEach(sstable -> {
pendingDirs.add(sstable.descriptor.directory);
});
if (pendingDir.exists())
for (File pendingDir : pendingDirs)
{
Preconditions.checkState(pendingDir.absolutePath().contains(transfer.planId.toString()));
logger.debug("Deleting pending transfer directory: {}", pendingDir);
pendingDir.deleteRecursive();
if (pendingDir.exists())
{
Preconditions.checkState(pendingDir.absolutePath().contains(transfer.planId.toString()));
logger.debug("Deleting pending transfer directory: {}", pendingDir);
pendingDir.deleteRecursive();
}
}
local.remove(transfer.planId);
@ -178,10 +145,8 @@ public class LocalTransfers
coordinating.remove(transfer.id());
CoordinatedTransfer.SingleTransferResult localPending = transfer.streamResults.get(FBUtilities.getBroadcastAddressAndPort());
PendingLocalTransfer localTransfer;
TimeUUID planId;
if (localPending != null && (planId = localPending.planId()) != null && (localTransfer = local.get(planId)) != null)
purge(localTransfer);
if (localPending != null)
purge(localPending.planId());
}
finally
{
@ -194,7 +159,7 @@ public class LocalTransfers
executor.submit(() -> {
try
{
cleanupCoordinatedTransfer(transfer);
purge(transfer);
}
catch (Throwable t)
{
@ -208,7 +173,7 @@ public class LocalTransfers
executor.submit(() -> {
try
{
cleanupPendingLocalTransfer(timeUUID);
purge(timeUUID);
}
catch (Throwable t)
{
@ -244,7 +209,7 @@ public class LocalTransfers
}
public static IVerbHandler<TransferFailed> verbHandler = message -> {
LocalTransfers.instance().purge(message.payload);
LocalTransfers.instance().purge(message.payload.planId);
MessagingService.instance().respond(NoPayload.noPayload, message);
};
}

View File

@ -80,14 +80,12 @@ 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)
{
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);
@ -95,8 +93,7 @@ public class PendingLocalTransfer
}
catch (Throwable t)
{
sstablesPriorToMove.forEach(s -> s.selfRef().release());
logger.error("Failed importing sstables");
moved.forEach(s -> s.selfRef().release());
for (Descriptor descriptor : movedDescriptors)
descriptor.getFormat().delete(descriptor);
throw new RuntimeException("Failed importing SSTables", t);

View File

@ -82,7 +82,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
CQLSSTableWriter.Builder builder1 = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + "(k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder1.build())
{
@ -93,7 +93,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
CQLSSTableWriter.Builder builder2 = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + "(k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder2.build())
@ -128,12 +128,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
Uninterruptibles.sleepUninterruptibly(3, TimeUnit.SECONDS);
// Assert that each node has 2 SSTables
cluster.forEach(instance -> {
instance.runOnInstance(() -> {
ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(KEYSPACE, "tbl");
assertEquals(2, cfs.getLiveSSTables().size());
});
});
assertSSTableCount(cluster, 2);
// Assert that each node has the correct values
assertLocalSelect(cluster, rows -> { assertRows(rows, row(1, 1), row(2, 1), row(3, 1)); });
@ -145,6 +140,10 @@ public class AccordImportSSTableTest extends TestBaseImpl
}
}
/**
* This might have a potential issue with us throwing an exception instead of having a custom class to notify
* that the read was a failure.
*/
@Test
public void testSSTableImportWithConcurrentTopologyChangeFails() throws Throwable
{
@ -153,13 +152,11 @@ public class AccordImportSSTableTest extends TestBaseImpl
CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + "(k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder.build())
{
writer.addRow(1, 1);
writer.addRow(2, 1);
writer.addRow(3, 1);
}
try (Cluster cluster = init(builder().withNodes(3)
@ -180,7 +177,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
Set<String> paths = Set.of(file);
Assertions.assertThatThrownBy(() -> cfs.importNewSSTables(paths, true, true, true, true, true, true, true))
.isInstanceOf(RuntimeException.class)
.hasMessageContaining("SSTable import failed because of a concurrent topology change");
.hasMessageContaining("Failed adding SSTables");
});
}, "importer");
importer.start();
@ -190,12 +187,11 @@ public class AccordImportSSTableTest extends TestBaseImpl
State.waitForTopologyChange.countDown();
});
importer.join();
Uninterruptibles.sleepUninterruptibly(3, TimeUnit.SECONDS);
cluster.forEach(instance -> instance.runOnInstance(() -> {
ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(KEYSPACE, TABLE);
assertEquals(0, cfs.getLiveSSTables().size());
}));
assertSSTableCount(cluster, 0);
}
}
@ -208,7 +204,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + "(k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder.build())
{
@ -256,7 +252,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + " (k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder.build())
{
@ -302,7 +298,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + " (k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder.build())
{
@ -348,7 +344,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + " (k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder.build())
{
@ -390,7 +386,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
// Wait until the recovery coordinator picks up the Import Txn
Uninterruptibles.sleepUninterruptibly(10, TimeUnit.SECONDS);
for (int i = 2; i <= 3; i++)
for (int i = 1; i <= 3; i++)
{
cluster.get(i).runOnInstance(() -> {
ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(KEYSPACE, TABLE);
@ -403,23 +399,111 @@ public class AccordImportSSTableTest extends TestBaseImpl
}
@Test
public void testImportOneTokenSSTable() throws Throwable
public void testRecoveryCoordinatorPerformsImport2() throws Throwable
{
String file = Files.createTempDirectory(AccordImportSSTableTest.class.getSimpleName()).toString();
CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + " (k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder.build())
{
writer.addRow(1, 1);
writer.addRow(2, 1);
writer.addRow(3, 1);
}
try (Cluster cluster = init(builder().withNodes(3).withoutVNodes()
.withDataDirCount(1).withConfig((config) ->
config
.set("accord.recover_txn", "100ms")
.set("accord.permit_local_delivery", false)
.with(Feature.NETWORK, Feature.GOSSIP)).start()))
{
cluster.schemaChange("DROP KEYSPACE IF EXISTS " + KEYSPACE);
cluster.schemaChange("CREATE KEYSPACE " + KEYSPACE + " WITH REPLICATION={'class':'SimpleStrategy', 'replication_factor': 3}");
cluster.schemaChange("CREATE TABLE " + KEYSPACE_TABLE + " (k int PRIMARY KEY, v int) WITH transactional_mode='full'");
// Node 1 sends the StableThenRead message to node 2 and then node 1 fails, so the
// only existence of the stable message is at node 2
cluster.filters().outbound().messagesMatching((from, to, msg) -> {
if (from == 1 && msg.verb() == Verb.ACCORD_STABLE_THEN_READ_REQ.id)
{
// We prevent nodes 1 & 3 from receiving the StableThenRead message,
// from node 1 and then prevent node 1 from receiving any more messages
cluster.filters().outbound().from(1).to(1, 3).drop();
cluster.filters().inbound().to(1).drop();
// We still want node 2 to receive the message, so the ImportTxn is
// stable, however once that is done we do not want to receive any more messages
// from node 1
if (to == 2)
{
cluster.filters().outbound().from(1).drop();
return false;
}
return true;
}
return false;
}).drop();
cluster.get(1).runOnInstance(() -> {
ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(KEYSPACE, TABLE);
Set<String> paths = Set.of(file);
Assertions.assertThatThrownBy(() -> cfs.importNewSSTables(paths, true, true, true, true, true, true, true));
});
// Wait until the recovery coordinator picks up the Import Txn
Uninterruptibles.sleepUninterruptibly(10, TimeUnit.SECONDS);
Iterable<IInvokableInstance> up = cluster.stream()
.filter(instance -> instance != cluster.get(1))
.collect(Collectors.toList());
assertSSTableCount(up, 1);
assertLocalSelect(up, rows -> { assertRows(rows, row(1, 1), row(2, 1), row(3, 1)); });
}
}
@Test
public void testImportSSTablesWithZeroCopyStreaming() throws Throwable
{
}
@Test
public void testImportSSTablesCleanupWithMultipleDataDirectories() throws Throwable
{
String file = Files.createTempDirectory(AccordImportSSTableTest.class.getSimpleName()).toString();
CQLSSTableWriter.Builder builder1 = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE_TABLE + "(k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder1.build())
{
writer.addRow(1, 1);
writer.addRow(2, 1);
}
CQLSSTableWriter.Builder builder2 = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE_TABLE + "(k, v) " + "VALUES (?, ?)");
try (CQLSSTableWriter writer = builder2.build())
{
writer.addRow(3, 1);
}
try (Cluster cluster = init(builder().withNodes(3)
.withoutVNodes()
.withDataDirCount(1)
.withDataDirCount(3)
.withConfig((config) ->
config
.with(Feature.NETWORK, Feature.GOSSIP)).start()))
@ -443,16 +527,12 @@ public class AccordImportSSTableTest extends TestBaseImpl
Uninterruptibles.sleepUninterruptibly(3, TimeUnit.SECONDS);
// Assert that each node has 1 SSTable
cluster.forEach(instance -> {
instance.runOnInstance(() -> {
ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(KEYSPACE, "tbl");
assertEquals(1, cfs.getLiveSSTables().size());
});
});
// Assert that each node has 2 SSTables
assertSSTableCount(cluster, 2);
// Assert that each node has the correct values
assertLocalSelect(cluster, rows -> { assertRows(rows, row(1, 1)); });
assertLocalSelect(cluster, rows -> { assertRows(rows, row(1, 1), row(2, 1), row(3, 1)); });
// Assert that SSTables are moved from the pending directories
assertPendingDirs(cluster, (File pendingUuidDir) -> {
@ -461,15 +541,10 @@ public class AccordImportSSTableTest extends TestBaseImpl
}
}
@Test
public void testWithZeroCopyStreaming() throws Throwable
{
}
/**
* There is a bug with this case, because we are going to be marked as stable and then perform the read which has the
* TxnImport logic contained within it. For ImportTxn's, we need to retry the read.
* TxnImport logic contained within it. For ImportTxn's we prevent this replica from making any progress in the case,
* where we can not properly execute a read import txn. We want to make sure it is marked as stale.
*/
@Test
public void testImportSSTableFailsActivation() throws Throwable
@ -479,7 +554,7 @@ public class AccordImportSSTableTest extends TestBaseImpl
CQLSSTableWriter.Builder builder = CQLSSTableWriter.builder()
.forTable(TABLE_SCHEMA_CQL)
.inDirectory(file)
.using("INSERT INTO " + KEYSPACE + ".tbl (k, v) " + "VALUES (?, ?)");
.using("INSERT INTO " + KEYSPACE_TABLE + " (k, v) " + "VALUES (?, ?)");
// We import a 1 token SSTable so the txn read logic is not run multiple of times by different CommandStores
try (CQLSSTableWriter writer = builder.build())
@ -506,14 +581,13 @@ public class AccordImportSSTableTest extends TestBaseImpl
Uninterruptibles.sleepUninterruptibly(30, TimeUnit.SECONDS);
cluster.forEach(instance -> {
instance.runOnInstance(() -> {
ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(KEYSPACE, "tbl");
assertEquals(1, cfs.getLiveSSTables().size());
});
});
Iterable<IInvokableInstance> up = cluster.stream()
.filter(instance -> instance != cluster.get(2))
.collect(Collectors.toList());
assertLocalSelect(cluster, rows -> { assertRows(rows, row(1, 1), row(2, 1), row(3, 1)); });
assertSSTableCount(up, 1);
assertLocalSelect(up, rows -> { assertRows(rows, row(1, 1)); });
}
}
@ -561,6 +635,17 @@ public class AccordImportSSTableTest extends TestBaseImpl
}
}
private static void assertSSTableCount(Iterable<IInvokableInstance> validate, int count)
{
for (IInvokableInstance instance : validate)
{
instance.runOnInstance(() -> {
ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(KEYSPACE, TABLE);
Assertions.assertThat(cfs.getLiveSSTables().size()).isEqualTo(count);
});
}
}
private static void assertLocalSelect(Iterable<IInvokableInstance> validate, IIsolatedExecutor.SerializableConsumer<Object[][]> onRows)
{
for (IInvokableInstance instance : validate)