diff --git a/CHANGES.txt b/CHANGES.txt index 1c0989def3..2878611652 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 5.1 + * Replace blocking wait with non-blocking delay in paxos repair (CASSANDRA-20983) * Implementation of CEP-55 - Generation of role names (CASSANDRA-20897) * Add cqlsh autocompletion for the identity mapping feature (CASSANDRA-20021) * Add DDL Guardrail enabling administrators to disallow creation/modification of keyspaces with durable_writes = false (CASSANDRA-20913) 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 7e5935f03d..62124a5a20 100644 --- a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupLocalCoordinator.java +++ b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupLocalCoordinator.java @@ -19,15 +19,18 @@ package org.apache.cassandra.service.paxos.cleanup; import java.util.Collection; +import java.util.Comparator; import java.util.Map; +import java.util.PriorityQueue; +import java.util.Queue; 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; +import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.db.DecoratedKey; @@ -58,6 +61,37 @@ public class PaxosCleanupLocalCoordinator extends AsyncFuture PRIORITY_COMPARATOR = new Comparator() + { + @Override + public int compare(DelayedRepair o1, DelayedRepair o2) + { + long delta = o1.scheduledAtNanos - o2.scheduledAtNanos; + if (delta > 0) + return 1; + if (delta < 0) + return -1; + return 0; + } + }; + + public DelayedRepair(UncommittedPaxosKey uncommitted, long sleepMillis) + { + this.uncommitted = uncommitted; + this.scheduledAtNanos = Clock.Global.nanoTime() + MILLISECONDS.toNanos(sleepMillis); + } + + public boolean isRunnable() + { + return Clock.Global.nanoTime() - scheduledAtNanos > 0; + } + } + private final UUID session; private final TableId tableId; private final TableMetadata table; @@ -69,6 +103,7 @@ public class PaxosCleanupLocalCoordinator extends AsyncFuture inflight = new ConcurrentHashMap<>(); + private final Queue delayed = new PriorityQueue<>(DelayedRepair.PRIORITY_COMPARATOR); private final PaxosTableRepairs tableRepairs; private PaxosCleanupLocalCoordinator(SharedContext ctx, UUID session, TableId tableId, Collection> ranges, CloseableIterator uncommittedIter, boolean autoRepair) @@ -125,11 +160,40 @@ public class PaxosCleanupLocalCoordinator extends AsyncFuture SECONDS.toMillis(1)) + logger.warn("Encountered ballot that is more than 1 second in the future, is there a clock sync issue? {}", uncommitted.ballot()); + + if (ballotElapsedMillis >= txnTimeoutMillis) + return false; + + long sleepMillis = txnTimeoutMillis - ballotElapsedMillis; + logger.info("Paxos auto repair encountered a potentially in progress ballot, sleeping {}ms to allow the in flight operation to finish", sleepMillis); + + delayed.add(new DelayedRepair(uncommitted, sleepMillis)); + ScheduledExecutors.scheduledFastTasks.schedule(this::scheduleKeyRepairsOrFinish, sleepMillis, MILLISECONDS); + + return true; + } + /** * Schedule as many key repairs as we can, up to the paralellism limit. If no repairs are scheduled and * none are in flight when the iterator is exhausted, the session will be finished */ - private void scheduleKeyRepairsOrFinish() + private synchronized void scheduleKeyRepairsOrFinish() { int parallelism = DatabaseDescriptor.getPaxosRepairParallelism(); Preconditions.checkArgument(parallelism > 0); @@ -141,18 +205,33 @@ public class PaxosCleanupLocalCoordinator extends AsyncFuture { if (result.wasSuccessful()) onKeyFinish(uncommitted.getKey()); @@ -175,24 +251,6 @@ 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))