From a6cd2727d02a88a10c910b4b62267aba82367efd Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Mon, 11 Jan 2010 20:14:53 +0000 Subject: [PATCH] inline IStage.executorService, removing the useless Stage wrappers git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@898049 13f79535-47bb-0310-9956-ffa450edef68 --- .../apache/cassandra/concurrent/IStage.java | 121 ------------------ .../concurrent/MultiThreadedStage.java | 95 -------------- .../concurrent/SingleThreadedStage.java | 102 --------------- .../cassandra/concurrent/StageManager.java | 42 ++++-- .../org/apache/cassandra/db/CommitLog.java | 4 +- .../cassandra/net/MessagingService.java | 2 +- .../cassandra/service/AntiEntropyService.java | 2 +- .../service/StorageLoadBalancer.java | 2 - .../cassandra/service/StorageProxy.java | 2 +- .../service/AntiEntropyServiceTest.java | 5 +- 10 files changed, 36 insertions(+), 341 deletions(-) delete mode 100644 src/java/org/apache/cassandra/concurrent/IStage.java delete mode 100644 src/java/org/apache/cassandra/concurrent/MultiThreadedStage.java delete mode 100644 src/java/org/apache/cassandra/concurrent/SingleThreadedStage.java diff --git a/src/java/org/apache/cassandra/concurrent/IStage.java b/src/java/org/apache/cassandra/concurrent/IStage.java deleted file mode 100644 index 69460d8b8d..0000000000 --- a/src/java/org/apache/cassandra/concurrent/IStage.java +++ /dev/null @@ -1,121 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.cassandra.concurrent; - -import java.util.concurrent.*; - -/** - * An abstraction for stages as described in the SEDA paper by Matt Welsh. - * For reference to the paper look over here - * SEDA: An Architecture for WellConditioned, - Scalable Internet Services. - */ - -public interface IStage -{ - /** - * Get the name of the associated stage. - * @return name of the associated stage. - */ - public String getName(); - - /** - * Get the thread pool used by this stage - * internally. - */ - public ExecutorService getInternalThreadPool(); - - /** - * This method is used to execute a piece of code on - * this stage. The idea is that the run() method - * of this Runnable instance is invoked on a thread from a - * thread pool that belongs to this stage. - * @param runnable instance whose run() method needs to be invoked. - */ - public void execute(Runnable runnable); - - /** - * This method is used to execute a piece of code on - * this stage which returns a Future pointer. The idea - * is that the call() method of this Runnable - * instance is invoked on a thread from a thread pool - * that belongs to this stage. - - * @param callable instance that needs to be invoked. - * @return the future return object from the callable. - */ - public Future execute(Callable callable); - - /** - * This method is used to submit tasks to this stage - * that execute periodically. - * - * @param command the task to execute. - * @param delay the time to delay first execution - * @param unit the time unit of the initialDelay and period parameters - * @return the future return object from the runnable. - */ - public ScheduledFuture schedule(Runnable command, long delay, TimeUnit unit); - - /** - * This method is used to submit tasks to this stage - * that execute periodically. - * @param command the task to execute. - * @param initialDelay the time to delay first execution - * @param period the period between successive executions - * @param unit the time unit of the initialDelay and period parameters - * @return the future return object from the runnable. - */ - public ScheduledFuture scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit); - - /** - * This method is used to submit tasks to this stage - * that execute periodically. - * @param command the task to execute. - * @param initialDelay the time to delay first execution - * @param delay the delay between the termination of one execution and the commencement of the next. - * @param unit the time unit of the initialDelay and delay parameters - * @return the future return object from the runnable. - */ - public ScheduledFuture scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit); - - /** - * Shutdown the stage. All the threads of this stage - * are forcefully shutdown. Any pending tasks on this - * stage could be dropped or the stage could wait for - * these tasks to be completed. This is however an - * implementation detail. - */ - public void shutdown(); - - /** - * Checks if the stage has been shutdown. - * @return true if shut down, otherwise false. - */ - public boolean isShutdown(); - - /** - * This method returns the number of tasks that are - * pending on this stage to be executed. - * @return task count. - */ - public long getPendingTasks(); - - public long getCompletedTasks(); -} diff --git a/src/java/org/apache/cassandra/concurrent/MultiThreadedStage.java b/src/java/org/apache/cassandra/concurrent/MultiThreadedStage.java deleted file mode 100644 index bf023cf1e9..0000000000 --- a/src/java/org/apache/cassandra/concurrent/MultiThreadedStage.java +++ /dev/null @@ -1,95 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - - -package org.apache.cassandra.concurrent; - -import java.util.concurrent.*; - -/** - * This class is an implementation of the IStage interface. In particular - * it is for a stage that has a thread pool with multiple threads. For details - * please refer to the IStage documentation. - */ - -public class MultiThreadedStage implements IStage -{ - private String name_; - private JMXEnabledThreadPoolExecutor executorService_; - - public MultiThreadedStage(String name, int numThreads) - { - name_ = name; - executorService_ = new JMXEnabledThreadPoolExecutor( numThreads, - numThreads, - Integer.MAX_VALUE, - TimeUnit.SECONDS, - new LinkedBlockingQueue(), - new NamedThreadFactory(name) - ); - } - - public String getName() - { - return name_; - } - - public ExecutorService getInternalThreadPool() - { - return executorService_; - } - - public Future execute(Callable callable) { - return executorService_.submit(callable); - } - - public void execute(Runnable runnable) { - executorService_.execute(runnable); - } - - public ScheduledFuture schedule(Runnable command, long delay, TimeUnit unit) - { - throw new UnsupportedOperationException("This operation is not supported"); - } - - public ScheduledFuture scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit) { - throw new UnsupportedOperationException("This operation is not supported"); - } - - public ScheduledFuture scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit) { - throw new UnsupportedOperationException("This operation is not supported"); - } - - public void shutdown() { - executorService_.shutdownNow(); - } - - public boolean isShutdown() - { - return executorService_.isShutdown(); - } - - public long getPendingTasks(){ - return executorService_.getPendingTasks(); - } - - public long getCompletedTasks() - { - return executorService_.getCompletedTasks(); - } -} diff --git a/src/java/org/apache/cassandra/concurrent/SingleThreadedStage.java b/src/java/org/apache/cassandra/concurrent/SingleThreadedStage.java deleted file mode 100644 index 87635b525c..0000000000 --- a/src/java/org/apache/cassandra/concurrent/SingleThreadedStage.java +++ /dev/null @@ -1,102 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.cassandra.concurrent; - -import java.util.concurrent.Callable; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Future; -import java.util.concurrent.ScheduledFuture; -import java.util.concurrent.TimeUnit; - -/** - * This class is an implementation of the IStage interface. In particular - * it is for a stage that has a thread pool with a single thread. For details - * please refer to the IStage documentation. - */ - -public class SingleThreadedStage implements IStage -{ - protected JMXEnabledThreadPoolExecutor executorService_; - private String name_; - - public SingleThreadedStage(String name) - { - executorService_ = new JMXEnabledThreadPoolExecutor(name); - name_ = name; - } - - /* Implementing the IStage interface methods */ - - public String getName() - { - return name_; - } - - public ExecutorService getInternalThreadPool() - { - return executorService_; - } - - public void execute(Runnable runnable) - { - executorService_.execute(runnable); - } - - public Future execute(Callable callable) - { - return executorService_.submit(callable); - } - - public ScheduledFuture schedule(Runnable command, long delay, TimeUnit unit) - { - //return executorService_.schedule(command, delay, unit); - throw new UnsupportedOperationException("This operation is not supported"); - } - - public ScheduledFuture scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit) - { - //return executorService_.scheduleAtFixedRate(command, initialDelay, period, unit); - throw new UnsupportedOperationException("This operation is not supported"); - } - - public ScheduledFuture scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit) - { - //return executorService_.scheduleWithFixedDelay(command, initialDelay, delay, unit); - throw new UnsupportedOperationException("This operation is not supported"); - } - - public void shutdown() - { - executorService_.shutdownNow(); - } - - public boolean isShutdown() - { - return executorService_.isShutdown(); - } - - public long getPendingTasks(){ - return executorService_.getPendingTasks(); - } - - public long getCompletedTasks() - { - return executorService_.getCompletedTasks(); - } -} diff --git a/src/java/org/apache/cassandra/concurrent/StageManager.java b/src/java/org/apache/cassandra/concurrent/StageManager.java index 4bf4b5bafe..75fc8add70 100644 --- a/src/java/org/apache/cassandra/concurrent/StageManager.java +++ b/src/java/org/apache/cassandra/concurrent/StageManager.java @@ -21,6 +21,10 @@ package org.apache.cassandra.concurrent; import java.util.HashMap; import java.util.Map; import java.util.Set; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; import org.apache.cassandra.net.MessagingService; @@ -35,7 +39,7 @@ import static org.apache.cassandra.config.DatabaseDescriptor.getConcurrentReader */ public class StageManager { - private static Map stageQueues = new HashMap(); + private static Map stages = new HashMap(); public final static String READ_STAGE = "ROW-READ-STAGE"; public final static String MUTATION_STAGE = "ROW-MUTATION-STAGE"; @@ -47,22 +51,33 @@ public class StageManager static { - stageQueues.put(MUTATION_STAGE, new MultiThreadedStage(MUTATION_STAGE, getConcurrentWriters())); - stageQueues.put(READ_STAGE, new MultiThreadedStage(READ_STAGE, getConcurrentReaders())); - stageQueues.put(STREAM_STAGE, new SingleThreadedStage(STREAM_STAGE)); - stageQueues.put(GOSSIP_STAGE, new SingleThreadedStage("GMFD")); - stageQueues.put(RESPONSE_STAGE, new MultiThreadedStage("RESPONSE-STAGE", MessagingService.MESSAGE_DESERIALIZE_THREADS)); - stageQueues.put(AE_SERVICE_STAGE, new SingleThreadedStage(AE_SERVICE_STAGE)); - stageQueues.put(LOADBALANCE_STAGE, new SingleThreadedStage(LOADBALANCE_STAGE)); + stages.put(MUTATION_STAGE, multiThreadedStage(MUTATION_STAGE, getConcurrentWriters())); + stages.put(READ_STAGE, multiThreadedStage(READ_STAGE, getConcurrentReaders())); + stages.put(RESPONSE_STAGE, multiThreadedStage("RESPONSE-STAGE", MessagingService.MESSAGE_DESERIALIZE_THREADS)); + // the rest are all single-threaded + stages.put(STREAM_STAGE, new JMXEnabledThreadPoolExecutor(STREAM_STAGE)); + stages.put(GOSSIP_STAGE, new JMXEnabledThreadPoolExecutor("GMFD")); + stages.put(AE_SERVICE_STAGE, new JMXEnabledThreadPoolExecutor(AE_SERVICE_STAGE)); + stages.put(LOADBALANCE_STAGE, new JMXEnabledThreadPoolExecutor(LOADBALANCE_STAGE)); + } + + private static ThreadPoolExecutor multiThreadedStage(String name, int numThreads) + { + return new JMXEnabledThreadPoolExecutor(numThreads, + numThreads, + Integer.MAX_VALUE, + TimeUnit.SECONDS, + new LinkedBlockingQueue(), + new NamedThreadFactory(name)); } /** * Retrieve a stage from the StageManager * @param stageName name of the stage to be retrieved. */ - public static IStage getStage(String stageName) + public static ThreadPoolExecutor getStage(String stageName) { - return stageQueues.get(stageName); + return stages.get(stageName); } /** @@ -70,11 +85,10 @@ public class StageManager */ public static void shutdown() { - Set stages = stageQueues.keySet(); - for ( String stage : stages ) + Set stages = StageManager.stages.keySet(); + for (String stage : stages) { - IStage registeredStage = stageQueues.get(stage); - registeredStage.shutdown(); + StageManager.stages.get(stage).shutdown(); } } } diff --git a/src/java/org/apache/cassandra/db/CommitLog.java b/src/java/org/apache/cassandra/db/CommitLog.java index 583d77835d..4a08d2ff94 100644 --- a/src/java/org/apache/cassandra/db/CommitLog.java +++ b/src/java/org/apache/cassandra/db/CommitLog.java @@ -280,7 +280,7 @@ public class CommitLog void recover(File[] clogs) throws IOException { Set tablesRecovered = new HashSet
(); - assert StageManager.getStage(StageManager.MUTATION_STAGE).getCompletedTasks() == 0; + assert StageManager.getStage(StageManager.MUTATION_STAGE).getCompletedTaskCount() == 0; int rows = 0; for (File file : clogs) { @@ -363,7 +363,7 @@ public class CommitLog } // wait for all the writes to finish on the mutation stage - while (StageManager.getStage(StageManager.MUTATION_STAGE).getCompletedTasks() < rows) + while (StageManager.getStage(StageManager.MUTATION_STAGE).getCompletedTaskCount() < rows) { try { diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index 8d73b9d9f0..6095befd83 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -493,7 +493,7 @@ public class MessagingService implements IFailureDetectionEventListener private static void enqueueRunnable(String stageName, Runnable runnable){ - IStage stage = StageManager.getStage(stageName); + ExecutorService stage = StageManager.getStage(stageName); if ( stage != null ) { diff --git a/src/java/org/apache/cassandra/service/AntiEntropyService.java b/src/java/org/apache/cassandra/service/AntiEntropyService.java index cfa6d459f5..a0ce096cdf 100644 --- a/src/java/org/apache/cassandra/service/AntiEntropyService.java +++ b/src/java/org/apache/cassandra/service/AntiEntropyService.java @@ -483,7 +483,7 @@ public class AntiEntropyService for (MerkleTree.RowHash minrow : minrows) range.addHash(minrow); - StageManager.getStage(StageManager.AE_SERVICE_STAGE).execute(this); + StageManager.getStage(StageManager.AE_SERVICE_STAGE).submit(this); logger.debug("Validated " + validated + " rows into AEService tree for " + cf); } diff --git a/src/java/org/apache/cassandra/service/StorageLoadBalancer.java b/src/java/org/apache/cassandra/service/StorageLoadBalancer.java index 1c70682849..52c04ae6a7 100644 --- a/src/java/org/apache/cassandra/service/StorageLoadBalancer.java +++ b/src/java/org/apache/cassandra/service/StorageLoadBalancer.java @@ -25,8 +25,6 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.apache.log4j.Logger; import org.apache.cassandra.concurrent.JMXEnabledThreadPoolExecutor; -import org.apache.cassandra.concurrent.SingleThreadedStage; -import org.apache.cassandra.concurrent.StageManager; import org.apache.cassandra.dht.Token; import org.apache.cassandra.gms.ApplicationState; import org.apache.cassandra.gms.EndPointState; diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 0c3bf5e142..2bf3c8c215 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -496,7 +496,7 @@ public class StorageProxy implements StorageProxyMBean for (ReadCommand command: commands) { Callable callable = new weakReadLocalCallable(command); - futures.add(StageManager.getStage(StageManager.READ_STAGE).execute(callable)); + futures.add(StageManager.getStage(StageManager.READ_STAGE).submit(callable)); } for (Future future : futures) { diff --git a/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java b/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java index 2abe571ab3..d3f31a1637 100644 --- a/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java +++ b/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java @@ -235,11 +235,12 @@ public class AntiEntropyServiceTest extends CleanupHelper Future flushAES() { - return StageManager.getStage(StageManager.AE_SERVICE_STAGE).execute(new Callable(){ + return StageManager.getStage(StageManager.AE_SERVICE_STAGE).submit(new Callable() + { public Boolean call() { return true; } - }); + }); } }