From e28556f71c938469cf8a37d700af7d2703c36657 Mon Sep 17 00:00:00 2001 From: Alan Wang Date: Fri, 31 Jul 2026 14:36:04 -0700 Subject: [PATCH] updates --- modules/accord | 2 +- .../service/accord/CoordinatedTransfer.java | 15 ++----------- .../service/accord/PendingLocalTransfer.java | 22 +------------------ .../cassandra/service/accord/txn/TxnRead.java | 8 +++++++ .../test/accord/AccordImportSSTableTest.java | 7 ++---- 5 files changed, 14 insertions(+), 40 deletions(-) diff --git a/modules/accord b/modules/accord index c12dd6f876..e2352bc34e 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit c12dd6f8767dccd0c7cd3480f938845c9e24cb94 +Subproject commit e2352bc34edaf9e8662068ba163519c7ba2a50fb diff --git a/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java b/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java index 7a41adfaa7..18371b53e2 100644 --- a/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java +++ b/src/java/org/apache/cassandra/service/accord/CoordinatedTransfer.java @@ -27,7 +27,6 @@ import java.util.Map; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; @@ -98,8 +97,6 @@ public class CoordinatedTransfer final Map nodeStreamingContext; SingleTransferResult streamResult = SingleTransferResult.Init(); - // If importTxnEpochMismatch is true, all replicas will deterministically not import the SSTables in the pending directory - volatile boolean importTxnEpochMismatch = false; public CoordinatedTransfer(UUID importID, TableMetadata tableMetadata, Map nodeStreamingContext, long streamingEpoch, TokenRange allSSTableRanges) { @@ -120,21 +117,13 @@ public class CoordinatedTransfer logger.debug("{} Executing Accord bulk transfer {}", logPrefix(), this); LocalTransfers.instance().save(this); stream(); - PendingLocalTransfer pendingLocalTransfer = LocalTransfers.instance().local.get(streamResult.planId); - CountDownLatch latch = new CountDownLatch(1); - pendingLocalTransfer.registerLatch(latch); try { performImportTxn(); - latch.await(); - if (importTxnEpochMismatch) - { - LocalTransfers.instance().scheduleCoordinatedTransferCleanup(this); - throw new RuntimeException("SSTable import failed because of a concurrent topology change; please retry the operation"); - } + } - catch (ReadTimeoutException | InterruptedException e) + catch (Exception e) { throw new RuntimeException("SSTable import failed locally; however the operation may still be applied by the recovery coordinator", e); } diff --git a/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java b/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java index 1fcf16ed05..c7f03fe996 100644 --- a/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java +++ b/src/java/org/apache/cassandra/service/accord/PendingLocalTransfer.java @@ -84,27 +84,7 @@ public class PendingLocalTransfer if (activated) return; - CoordinatedTransfer coordinatedTransfer = LocalTransfers.instance.coordinating.get(metadata.getImportID()); - boolean isCoordinator = coordinatedTransfer != null; - - if (metadata.getStreamingEpoch() != executeAtEpoch) - { - logger.info("{} Failing activation of pending SSTables because streaming epoch {} != importTxn executeAt epoch {}", - logPrefix(), metadata.getStreamingEpoch(), executeAtEpoch); - - if (isCoordinator) - { - Invariants.require(latch != null); - latch.countDown(); - coordinatedTransfer.importTxnEpochMismatch = true; - } - - LocalTransfers.instance().schedulePendingLocalTransferCleanup(planId); - - activated = true; - return; - } - + Invariants.require(metadata.getStreamingEpoch() == executeAtEpoch); long startedActivation = currentTimeMillis(); logger.info("{} Activating transfer {}, {} ms since pending", logPrefix(), this, startedActivation - createdAt); ColumnFamilyStore cfs = ColumnFamilyStore.getIfExists(tableId); diff --git a/src/java/org/apache/cassandra/service/accord/txn/TxnRead.java b/src/java/org/apache/cassandra/service/accord/txn/TxnRead.java index a66dbc1adb..6537612097 100644 --- a/src/java/org/apache/cassandra/service/accord/txn/TxnRead.java +++ b/src/java/org/apache/cassandra/service/accord/txn/TxnRead.java @@ -535,6 +535,14 @@ public class TxnRead extends AbstractKeySorted implements Read return importMetadata != null; } + @Override + public long getImportStreamingEpoch() + { + if (importMetadata != null) + return importMetadata.streamingEpoch; + return -1L; + } + public static final ParameterisedVersionedSerializer serializer = new ParameterisedVersionedSerializer<>() { @Override 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 5b1c7f28e9..580a80f782 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordImportSSTableTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordImportSSTableTest.java @@ -170,10 +170,7 @@ public class AccordImportSSTableTest extends TestBaseImpl Assertions.assertThatThrownBy(() -> cfs.importNewSSTables(paths, true, true, true, true, true, true, true)) .isInstanceOf(RuntimeException.class) .hasMessageContaining("Failed adding SSTables on local node; note the import may still have been committed by a recovery coordinator") - .cause() - .isInstanceOf(RuntimeException.class) - .hasMessageContaining("SSTable import failed because of a concurrent topology change; please retry the operation"); - + .cause(); }); }, "importer"); @@ -453,7 +450,7 @@ public class AccordImportSSTableTest extends TestBaseImpl cfs.importNewSSTables(Set.of(file), true, true, true, true, true, true, true); }); - Uninterruptibles.sleepUninterruptibly(10, TimeUnit.SECONDS); + Uninterruptibles.sleepUninterruptibly(15, TimeUnit.SECONDS); assertLocalSelect(cluster, rows -> assertRows(rows, row(1, 1), row(2, 1), row(3, 1))); assertSSTableCount(cluster, 1);