From a449a4f76baf41b4707f177df208131589f981bb Mon Sep 17 00:00:00 2001 From: Ariel Weisberg Date: Fri, 21 Mar 2025 15:45:17 -0400 Subject: [PATCH] PaxosCleanupLocalCoordinator wait for transaction timeout before repairing Patch by Ariel Weisberg; Reviewed by Benedict Elliott Smith for CASSANDRA-20469 --- CHANGES.txt | 1 + .../org/apache/cassandra/config/Config.java | 2 ++ .../cassandra/config/DatabaseDescriptor.java | 20 ++++++++--- .../cassandra/service/StorageService.java | 14 +++++++- .../service/StorageServiceMBean.java | 4 +++ .../cleanup/PaxosCleanupLocalCoordinator.java | 34 +++++++++++++++++-- 6 files changed, 67 insertions(+), 8 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index f4e934ab16..04a146aa02 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 5.1 + * Fix Paxos repair interrupts running transactions (CASSANDRA-20469) * Various fixes in constraint framework (CASSANDRA-20481) * Add support in CAS for -= on numeric types, and fixed improper handling of empty bytes which lead to NPE (CASSANDRA-20477) * Do not fail to start a node with materialized views after they are turned off in config (CASSANDRA-20452) diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index 2ec1d78e30..2943c973ed 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -1467,4 +1467,6 @@ public class Config // 3.x Cassandra Driver has its "read" timeout set to 12 seconds, default matches this. public DurationSpec.LongMillisecondsBound native_transport_timeout = new DurationSpec.LongMillisecondsBound("12s"); public boolean enforce_native_deadline_for_hints = false; + + public boolean paxos_repair_race_wait = true; } diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index f6fd1b52ff..6afcabf5fb 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -61,7 +61,6 @@ import com.google.common.collect.Sets; import com.google.common.primitives.Ints; import com.google.common.primitives.Longs; import com.google.common.util.concurrent.RateLimiter; - import org.apache.commons.lang3.ArrayUtils; import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; @@ -105,13 +104,13 @@ import org.apache.cassandra.locator.DynamicEndpointSnitch; import org.apache.cassandra.locator.EndpointSnitchInfo; import org.apache.cassandra.locator.IEndpointSnitch; import org.apache.cassandra.locator.InetAddressAndPort; -import org.apache.cassandra.locator.Locator; -import org.apache.cassandra.locator.LocationInfo; import org.apache.cassandra.locator.InitialLocationProvider; +import org.apache.cassandra.locator.LocationInfo; +import org.apache.cassandra.locator.Locator; import org.apache.cassandra.locator.NodeAddressConfig; +import org.apache.cassandra.locator.NodeProximity; import org.apache.cassandra.locator.ReconnectableSnitchHelper; import org.apache.cassandra.locator.SeedProvider; -import org.apache.cassandra.locator.NodeProximity; import org.apache.cassandra.locator.SnitchAdapter; import org.apache.cassandra.security.AbstractCryptoProvider; import org.apache.cassandra.security.EncryptionContext; @@ -129,8 +128,8 @@ import org.apache.cassandra.utils.StorageCompatibilityMode; import static org.apache.cassandra.config.CassandraRelevantProperties.ALLOCATE_TOKENS_FOR_KEYSPACE; import static org.apache.cassandra.config.CassandraRelevantProperties.ALLOW_UNLIMITED_CONCURRENT_VALIDATIONS; import static org.apache.cassandra.config.CassandraRelevantProperties.AUTO_BOOTSTRAP; -import static org.apache.cassandra.config.CassandraRelevantProperties.CONFIG_LOADER; import static org.apache.cassandra.config.CassandraRelevantProperties.CHRONICLE_ANALYTICS_DISABLE; +import static org.apache.cassandra.config.CassandraRelevantProperties.CONFIG_LOADER; import static org.apache.cassandra.config.CassandraRelevantProperties.DISABLE_STCS_IN_L0; import static org.apache.cassandra.config.CassandraRelevantProperties.INITIAL_TOKEN; import static org.apache.cassandra.config.CassandraRelevantProperties.IO_NETTY_TRANSPORT_ESTIMATE_SIZE_ON_SUBMIT; @@ -5574,4 +5573,15 @@ public class DatabaseDescriptor { conf.tombstone_read_purgeable_metric_granularity = granularity; } + + public static boolean getPaxosRepairRaceWait() + { + return conf.paxos_repair_race_wait; + } + + @VisibleForTesting + public static void setPaxosRepairRaceWait(boolean paxosRepairRaceWait) + { + conf.paxos_repair_race_wait = paxosRepairRaceWait; + } } diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 812dc31cd4..e8a3b40a86 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -142,11 +142,11 @@ import org.apache.cassandra.locator.EndpointsForToken; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.LocalStrategy; import org.apache.cassandra.locator.MetaStrategy; +import org.apache.cassandra.locator.NodeProximity; import org.apache.cassandra.locator.RangesAtEndpoint; import org.apache.cassandra.locator.RangesByEndpoint; import org.apache.cassandra.locator.Replica; import org.apache.cassandra.locator.Replicas; -import org.apache.cassandra.locator.NodeProximity; import org.apache.cassandra.locator.SnitchAdapter; import org.apache.cassandra.locator.SystemReplicas; import org.apache.cassandra.metrics.Sampler; @@ -5541,4 +5541,16 @@ public class StorageService extends NotificationBroadcasterSupport implements IE { DatabaseDescriptor.setPrioritizeSAIOverLegacyIndex(value); } + + @Override + public void setPaxosRepairRaceWait(boolean paxosRepairRaceWait) + { + DatabaseDescriptor.setPaxosRepairRaceWait(paxosRepairRaceWait); + } + + @Override + public boolean getPaxosRepairRaceWait() + { + return DatabaseDescriptor.getPaxosRepairRaceWait(); + } } diff --git a/src/java/org/apache/cassandra/service/StorageServiceMBean.java b/src/java/org/apache/cassandra/service/StorageServiceMBean.java index 7760ea0188..b738ecd486 100644 --- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java +++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java @@ -1356,4 +1356,8 @@ public interface StorageServiceMBean extends NotificationEmitter boolean getPrioritizeSAIOverLegacyIndex(); void setPrioritizeSAIOverLegacyIndex(boolean value); + + void setPaxosRepairRaceWait(boolean paxosRepairCoordinatorWait); + + boolean getPaxosRepairRaceWait(); } diff --git a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupLocalCoordinator.java b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupLocalCoordinator.java index a53fec3e60..7e5935f03d 100644 --- a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupLocalCoordinator.java +++ b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupLocalCoordinator.java @@ -24,6 +24,7 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import com.google.common.base.Preconditions; +import com.google.common.util.concurrent.Uninterruptibles; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -41,9 +42,15 @@ import org.apache.cassandra.service.paxos.PaxosRepair; import org.apache.cassandra.service.paxos.PaxosState; import org.apache.cassandra.service.paxos.uncommitted.UncommittedPaxosKey; import org.apache.cassandra.tcm.ClusterMetadata; +import org.apache.cassandra.utils.Clock; import org.apache.cassandra.utils.CloseableIterator; import org.apache.cassandra.utils.concurrent.AsyncFuture; +import static java.util.concurrent.TimeUnit.MICROSECONDS; +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.apache.cassandra.config.DatabaseDescriptor.getCasContentionTimeout; +import static org.apache.cassandra.config.DatabaseDescriptor.getWriteRpcTimeout; import static org.apache.cassandra.service.paxos.cleanup.PaxosCleanupSession.TIMEOUT_NANOS; public class PaxosCleanupLocalCoordinator extends AsyncFuture @@ -134,8 +141,10 @@ public class PaxosCleanupLocalCoordinator extends AsyncFuture { if (result.wasSuccessful()) onKeyFinish(uncommitted.getKey()); @@ -163,6 +175,24 @@ public class PaxosCleanupLocalCoordinator extends AsyncFuture SECONDS.toMicros(1)) + logger.warn("Encountered ballot that is more than 1 second in the future, is there a clock sync issue? {}", uncommitted.ballot()); + if (ballotElapsedMicros < txnTimeoutMicros) + { + long sleepMicros = txnTimeoutMicros - ballotElapsedMicros; + logger.info("Paxos auto repair encountered a potentially in progress ballot, sleeping {}us to allow the in flight operation to finish", sleepMicros); + Uninterruptibles.sleepUninterruptibly(sleepMicros, MICROSECONDS); + } + } + private synchronized void onKeyFinish(DecoratedKey key) { if (!inflight.containsKey(key))