diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index b841d9b7c0..21ca1b595c 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -1264,4 +1264,6 @@ public class Config // 3.0 Cassandra Driver has its "read" timeout set to 12 seconds. Our recommendation is match this. public DurationSpec.LongMillisecondsBound native_transport_timeout = new DurationSpec.LongMillisecondsBound("12000ms"); 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 853b29f6df..55556a6308 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -4638,4 +4638,15 @@ public class DatabaseDescriptor { conf.reject_out_of_token_range_requests = enabled; } + + 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 6de663910c..8d494fc4d6 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -7153,4 +7153,16 @@ public class StorageService extends NotificationBroadcasterSupport implements IE DatabaseDescriptor.setEnforceNativeDeadlineForHints(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 c43ef2eec7..c7c02d2750 100644 --- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java +++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java @@ -1103,4 +1103,8 @@ public interface StorageServiceMBean extends NotificationEmitter * e.g. keyspace_name -> [reads, writes, paxos]. */ Map getOutOfRangeOperationCounts(); + + 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 3904d54d32..f7f500f57e 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; @@ -39,9 +40,15 @@ import org.apache.cassandra.service.paxos.AbstractPaxosRepair; 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.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; import static org.apache.cassandra.utils.Clock.Global.nanoTime; @@ -126,8 +133,10 @@ public class PaxosCleanupLocalCoordinator extends AsyncFuture { if (result.wasSuccessful()) onKeyFinish(uncommitted.getKey()); @@ -155,6 +167,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))