From 1b914f1a8e9e0b826fbd33f586ecde80655afd5d Mon Sep 17 00:00:00 2001 From: Alan Wang Date: Mon, 20 Jul 2026 13:59:06 -0700 Subject: [PATCH] update --- .../service/accord/CoordinatedTransfer.java | 2 +- .../service/accord/LocalTransfers.java | 81 +++----- .../service/accord/PendingLocalTransfer.java | 5 +- .../test/accord/AccordImportSSTableTest.java | 179 +++++++++++++----- 4 files changed, 157 insertions(+), 110 deletions(-) diff --git a/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java b/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java index 2702e70b4f..cea7386183 100644 --- a/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java +++ b/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java @@ -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 diff --git a/src/java/org/apache/cassandra/service/accord/LocalTransfers.java b/src/java/org/apache/cassandra/service/accord/LocalTransfers.java index 9ff16fd3f6..da4eed260b 100644 --- a/src/java/org/apache/cassandra/service/accord/LocalTransfers.java +++ b/src/java/org/apache/cassandra/service/accord/LocalTransfers.java @@ -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// if (!transfer.sstables.isEmpty()) { - SSTableReader sstable = transfer.sstables.iterator().next(); - File pendingDir = sstable.descriptor.directory; + Set 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 verbHandler = message -> { - LocalTransfers.instance().purge(message.payload); + LocalTransfers.instance().purge(message.payload.planId); MessagingService.instance().respond(NoPayload.noPayload, message); }; } diff --git a/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java b/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java index 3cb35a32cc..71c96c8fa9 100644 --- a/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java +++ b/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java @@ -80,14 +80,12 @@ 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) { 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); diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordImportSSTableTest.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordImportSSTableTest.java index e4fb72b1d2..db398830c8 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordImportSSTableTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordImportSSTableTest.java @@ -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 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 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 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 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 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 validate, IIsolatedExecutor.SerializableConsumer onRows) { for (IInvokableInstance instance : validate)