From 3b970ddbd6c0018ee13b596dddfd838c43afa4d1 Mon Sep 17 00:00:00 2001 From: Benjamin Lerer Date: Tue, 10 Nov 2020 09:55:20 +0100 Subject: [PATCH] Remove the SEPExecutor blocking behavior MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit patch by Benjamin Lerer; reviewed by Andrés de la Peña for CASSANDRA-16186 If the number of tasks added to a SEPExecutor exceed the max queue size, the threads adding those task will be block until enough space become available for all the blocked tasks. At this point all the blocked threads will be released at once. As the maxTasksQueued was always set to INTEGER.MAX_VALUE, the code was in practice never using. By consequence, removing the code was a better option than trying to fix it. --- CHANGES.txt | 1 + .../cassandra/concurrent/SEPExecutor.java | 34 +------------------ .../concurrent/SharedExecutorPool.java | 4 +-- .../cassandra/concurrent/StageManager.java | 2 +- .../apache/cassandra/transport/Message.java | 1 - .../LongSharedExecutorPoolTest.java | 8 ++--- .../cassandra/concurrent/SEPExecutorTest.java | 2 +- 7 files changed, 10 insertions(+), 42 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 07ee9af879..35612cf5ba 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 3.0.24: + * Remove the SEPExecutor blocking behavior (CASSANDRA-16186) * Wait for schema agreement when bootstrapping (CASSANDRA-15158) * Fix invalid cell value skipping when reading from disk (CASSANDRA-16223) * Prevent invoking enable/disable gossip when not in NORMAL (CASSANDRA-16146) diff --git a/src/java/org/apache/cassandra/concurrent/SEPExecutor.java b/src/java/org/apache/cassandra/concurrent/SEPExecutor.java index d5c7b14219..190004d426 100644 --- a/src/java/org/apache/cassandra/concurrent/SEPExecutor.java +++ b/src/java/org/apache/cassandra/concurrent/SEPExecutor.java @@ -25,7 +25,6 @@ import java.util.concurrent.atomic.AtomicLong; import org.apache.cassandra.metrics.SEPMetrics; import org.apache.cassandra.utils.concurrent.SimpleCondition; -import org.apache.cassandra.utils.concurrent.WaitQueue; import static org.apache.cassandra.concurrent.SEPWorker.Work; @@ -34,7 +33,6 @@ public class SEPExecutor extends AbstractLocalAwareExecutorService private final SharedExecutorPool pool; public final int maxWorkers; - private final int maxTasksQueued; private final SEPMetrics metrics; // stores both a set of work permits and task permits: @@ -42,8 +40,6 @@ public class SEPExecutor extends AbstractLocalAwareExecutorService // top 32 bits are number of work permits available in the range [0..maxWorkers] (initially maxWorkers) private final AtomicLong permits = new AtomicLong(); - // producers wait on this when there is no room on the queue - private final WaitQueue hasRoom = new WaitQueue(); private final AtomicLong completedTasks = new AtomicLong(); volatile boolean shuttingDown = false; @@ -52,11 +48,10 @@ public class SEPExecutor extends AbstractLocalAwareExecutorService // TODO: see if other queue implementations might improve throughput protected final ConcurrentLinkedQueue> tasks = new ConcurrentLinkedQueue<>(); - SEPExecutor(SharedExecutorPool pool, int maxWorkers, int maxTasksQueued, String jmxPath, String name) + SEPExecutor(SharedExecutorPool pool, int maxWorkers, String jmxPath, String name) { this.pool = pool; this.maxWorkers = maxWorkers; - this.maxTasksQueued = maxTasksQueued; this.permits.set(combine(0, maxWorkers)); this.metrics = new SEPMetrics(this, jmxPath, name); } @@ -102,29 +97,6 @@ public class SEPExecutor extends AbstractLocalAwareExecutorService // worker, we simply start a worker in a spinning state pool.maybeStartSpinningWorker(); } - else if (taskPermits >= maxTasksQueued) - { - // register to receive a signal once a task is processed bringing the queue below its threshold - WaitQueue.Signal s = hasRoom.register(); - - // we will only be signalled once the queue drops below full, so this creates equivalent external behaviour - // however the advantage is that we never wake-up spuriously; - // we choose to always sleep, even if in the intervening time the queue has dropped below limit, - // so long as we _will_ eventually receive a signal - if (taskPermits(permits.get()) > maxTasksQueued) - { - // if we're blocking, we might as well directly schedule a worker if we aren't already at max - if (takeWorkPermit(true)) - pool.schedule(new Work(this)); - - metrics.totalBlocked.inc(); - metrics.currentBlocked.inc(); - s.awaitUninterruptibly(); - metrics.currentBlocked.dec(); - } - else // don't propagate our signal when we cancel, just cancel - s.cancel(); - } } // takes permission to perform a task, if any are available; once taken it is guaranteed @@ -139,8 +111,6 @@ public class SEPExecutor extends AbstractLocalAwareExecutorService return false; if (permits.compareAndSet(current, updateTaskPermits(current, taskPermits - 1))) { - if (taskPermits == maxTasksQueued && hasRoom.hasWaiters()) - hasRoom.signalAll(); return true; } } @@ -159,8 +129,6 @@ public class SEPExecutor extends AbstractLocalAwareExecutorService return false; if (permits.compareAndSet(current, combine(taskPermits - taskDelta, workPermits - 1))) { - if (takeTaskPermit && taskPermits == maxTasksQueued && hasRoom.hasWaiters()) - hasRoom.signalAll(); return true; } } diff --git a/src/java/org/apache/cassandra/concurrent/SharedExecutorPool.java b/src/java/org/apache/cassandra/concurrent/SharedExecutorPool.java index e6a0df784a..ff21771a74 100644 --- a/src/java/org/apache/cassandra/concurrent/SharedExecutorPool.java +++ b/src/java/org/apache/cassandra/concurrent/SharedExecutorPool.java @@ -108,9 +108,9 @@ public class SharedExecutorPool schedule(Work.SPINNING); } - public synchronized LocalAwareExecutorService newExecutor(int maxConcurrency, int maxQueuedTasks, String jmxPath, String name) + public synchronized LocalAwareExecutorService newExecutor(int maxConcurrency, String jmxPath, String name) { - SEPExecutor executor = new SEPExecutor(this, maxConcurrency, maxQueuedTasks, jmxPath, name); + SEPExecutor executor = new SEPExecutor(this, maxConcurrency, jmxPath, name); executors.add(executor); return executor; } diff --git a/src/java/org/apache/cassandra/concurrent/StageManager.java b/src/java/org/apache/cassandra/concurrent/StageManager.java index 457e8011ed..b609c6cef8 100644 --- a/src/java/org/apache/cassandra/concurrent/StageManager.java +++ b/src/java/org/apache/cassandra/concurrent/StageManager.java @@ -91,7 +91,7 @@ public class StageManager private static LocalAwareExecutorService multiThreadedLowSignalStage(Stage stage, int numThreads) { - return SharedExecutorPool.SHARED.newExecutor(numThreads, Integer.MAX_VALUE, stage.getJmxType(), stage.getJmxName()); + return SharedExecutorPool.SHARED.newExecutor(numThreads, stage.getJmxType(), stage.getJmxName()); } /** diff --git a/src/java/org/apache/cassandra/transport/Message.java b/src/java/org/apache/cassandra/transport/Message.java index 52025782b1..a88540ba56 100644 --- a/src/java/org/apache/cassandra/transport/Message.java +++ b/src/java/org/apache/cassandra/transport/Message.java @@ -421,7 +421,6 @@ public abstract class Message public static class Dispatcher extends SimpleChannelInboundHandler { private static final LocalAwareExecutorService requestExecutor = SHARED.newExecutor(DatabaseDescriptor.getNativeTransportMaxThreads(), - Integer.MAX_VALUE, "transport", "Native-Transport-Requests"); diff --git a/test/burn/org/apache/cassandra/concurrent/LongSharedExecutorPoolTest.java b/test/burn/org/apache/cassandra/concurrent/LongSharedExecutorPoolTest.java index fe464c7804..f6c572081d 100644 --- a/test/burn/org/apache/cassandra/concurrent/LongSharedExecutorPoolTest.java +++ b/test/burn/org/apache/cassandra/concurrent/LongSharedExecutorPoolTest.java @@ -106,7 +106,7 @@ public class LongSharedExecutorPoolTest { final int executorCount = 4; int threadCount = 8; - int maxQueued = 1024; + int scale = 1024; final WeibullDistribution workTime = new WeibullDistribution(3, 200000); final long minWorkTime = TimeUnit.MICROSECONDS.toNanos(1); final long maxWorkTime = TimeUnit.MILLISECONDS.toNanos(1); @@ -116,11 +116,11 @@ public class LongSharedExecutorPoolTest final ExecutorService[] executors = new ExecutorService[executorCount]; for (int i = 0 ; i < executors.length ; i++) { - executors[i] = SharedExecutorPool.SHARED.newExecutor(threadCount, maxQueued, "test" + i, "test" + i); + executors[i] = SharedExecutorPool.SHARED.newExecutor(threadCount, "test" + i, "test" + i); threadCounts[i] = threadCount; - workCount[i] = new WeibullDistribution(2, maxQueued); + workCount[i] = new WeibullDistribution(2, scale); threadCount *= 2; - maxQueued *= 2; + scale *= 2; } long runs = 0; diff --git a/test/unit/org/apache/cassandra/concurrent/SEPExecutorTest.java b/test/unit/org/apache/cassandra/concurrent/SEPExecutorTest.java index 8526dd04fd..9e98cf1294 100644 --- a/test/unit/org/apache/cassandra/concurrent/SEPExecutorTest.java +++ b/test/unit/org/apache/cassandra/concurrent/SEPExecutorTest.java @@ -51,7 +51,7 @@ public class SEPExecutorTest for (int idx = 0; idx < 20; idx++) { - ExecutorService es = sharedPool.newExecutor(FBUtilities.getAvailableProcessors(), Integer.MAX_VALUE, "STAGE", run + MAGIC + idx); + ExecutorService es = sharedPool.newExecutor(FBUtilities.getAvailableProcessors(), "STAGE", run + MAGIC + idx); // Write to black hole es.execute(() -> nullPrintSteam.println("TEST" + es)); }