From 6b5d8bf8028c04db5d167543d006da901eb01663 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Mon, 11 Jan 2010 19:59:58 +0000 Subject: [PATCH] centralize stage creation in StageManager and standardize stage naming patch by jbellis; reviewed by gdusbabek for CASSANDRA-684 git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@898035 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/concurrent/StageManager.java | 95 +++++-------------- .../org/apache/cassandra/db/CommitLog.java | 6 +- .../org/apache/cassandra/db/RangeCommand.java | 2 +- .../cassandra/db/RangeSliceCommand.java | 2 +- .../org/apache/cassandra/db/ReadCommand.java | 2 +- .../org/apache/cassandra/db/ReadResponse.java | 3 +- .../org/apache/cassandra/db/RowMutation.java | 2 +- .../cassandra/db/RowMutationMessage.java | 2 +- .../org/apache/cassandra/gms/Gossiper.java | 11 +-- .../cassandra/io/StreamRequestMessage.java | 2 +- .../org/apache/cassandra/net/Message.java | 3 +- .../cassandra/net/MessagingService.java | 23 ++--- .../cassandra/net/TcpConnectionManager.java | 6 +- .../cassandra/service/AntiEntropyService.java | 12 +-- .../service/StorageLoadBalancer.java | 2 - .../cassandra/service/StorageProxy.java | 7 +- .../cassandra/service/StorageService.java | 2 +- .../service/AntiEntropyServiceTest.java | 2 +- 18 files changed, 60 insertions(+), 124 deletions(-) diff --git a/src/java/org/apache/cassandra/concurrent/StageManager.java b/src/java/org/apache/cassandra/concurrent/StageManager.java index f1500f3528..4bf4b5bafe 100644 --- a/src/java/org/apache/cassandra/concurrent/StageManager.java +++ b/src/java/org/apache/cassandra/concurrent/StageManager.java @@ -21,55 +21,39 @@ package org.apache.cassandra.concurrent; import java.util.HashMap; import java.util.Map; import java.util.Set; -import java.util.concurrent.ExecutorService; + +import org.apache.cassandra.net.MessagingService; import static org.apache.cassandra.config.DatabaseDescriptor.getConcurrentWriters; import static org.apache.cassandra.config.DatabaseDescriptor.getConcurrentReaders; /** - * This class manages all stages that exist within a process. The application registers - * and de-registers stages with this abstraction. Any component that has the ID - * associated with a stage can obtain a handle to actual stage. + * This class manages executor services for Messages recieved: each Message requests + * running on a specific "stage" for concurrency control; hence the Map approach, + * even though stages (executors) are not created dynamically. */ - public class StageManager { - private static Map stageQueues_ = new HashMap(); + private static Map stageQueues = new HashMap(); - public final static String readStage_ = "ROW-READ-STAGE"; - public final static String mutationStage_ = "ROW-MUTATION-STAGE"; - public final static String streamStage_ = "STREAM-STAGE"; + public final static String READ_STAGE = "ROW-READ-STAGE"; + public final static String MUTATION_STAGE = "ROW-MUTATION-STAGE"; + public final static String STREAM_STAGE = "STREAM-STAGE"; + public final static String GOSSIP_STAGE = "GS"; + public static final String RESPONSE_STAGE = "RESPONSE-STAGE"; + public final static String AE_SERVICE_STAGE = "AE-SERVICE-STAGE"; + private static final String LOADBALANCE_STAGE = "LOAD-BALANCER-STAGE"; static { - StageManager.registerStage(mutationStage_, new MultiThreadedStage(mutationStage_, getConcurrentWriters())); - StageManager.registerStage(readStage_, new MultiThreadedStage(readStage_, getConcurrentReaders())); - StageManager.registerStage(streamStage_, new SingleThreadedStage(streamStage_)); - } - - /** - * Register a stage with the StageManager - * @param stageName stage name. - * @param stage stage for the respective message types. - */ - public static void registerStage(String stageName, IStage stage) - { - stageQueues_.put(stageName, stage); - } - - /** - * Returns the stage that we are currently executing on. - * This relies on the fact that the thread names in the - * stage have the name of the stage as the prefix. - * @return Returns the stage that we are currently executing on. - */ - public static IStage getCurrentStage() - { - String name = Thread.currentThread().getName(); - String[] peices = name.split(":"); - IStage stage = getStage(peices[0]); - return stage; + 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)); } /** @@ -78,51 +62,18 @@ public class StageManager */ public static IStage getStage(String stageName) { - return stageQueues_.get(stageName); + return stageQueues.get(stageName); } - /** - * Retrieve the internal thread pool associated with the - * specified stage name. - * @param stageName name of the stage. - */ - public static ExecutorService getStageInternalThreadPool(String stageName) - { - IStage stage = getStage(stageName); - if ( stage == null ) - throw new IllegalArgumentException("No stage registered with name " + stageName); - return stage.getInternalThreadPool(); - } - - /** - * Deregister a stage from StageManager - * @param stageName stage name. - */ - public static void deregisterStage(String stageName) - { - stageQueues_.remove(stageName); - } - - /** - * This method gets the number of tasks on the - * stage's internal queue. - * @param stage name of the stage - * @return stage task count. - */ - public static long getStageTaskCount(String stage) - { - return stageQueues_.get(stage).getPendingTasks(); - } - /** * This method shuts down all registered stages. */ public static void shutdown() { - Set stages = stageQueues_.keySet(); + Set stages = stageQueues.keySet(); for ( String stage : stages ) { - IStage registeredStage = stageQueues_.get(stage); + IStage registeredStage = stageQueues.get(stage); registeredStage.shutdown(); } } diff --git a/src/java/org/apache/cassandra/db/CommitLog.java b/src/java/org/apache/cassandra/db/CommitLog.java index 8c4ae29702..583d77835d 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.mutationStage_).getCompletedTasks() == 0; + assert StageManager.getStage(StageManager.MUTATION_STAGE).getCompletedTasks() == 0; int rows = 0; for (File file : clogs) { @@ -356,14 +356,14 @@ public class CommitLog } } }; - StageManager.getStage(StageManager.mutationStage_).execute(runnable); + StageManager.getStage(StageManager.MUTATION_STAGE).execute(runnable); rows++; } reader.close(); } // wait for all the writes to finish on the mutation stage - while (StageManager.getStage(StageManager.mutationStage_).getCompletedTasks() < rows) + while (StageManager.getStage(StageManager.MUTATION_STAGE).getCompletedTasks() < rows) { try { diff --git a/src/java/org/apache/cassandra/db/RangeCommand.java b/src/java/org/apache/cassandra/db/RangeCommand.java index 829052b829..c4ff65c185 100644 --- a/src/java/org/apache/cassandra/db/RangeCommand.java +++ b/src/java/org/apache/cassandra/db/RangeCommand.java @@ -55,7 +55,7 @@ public class RangeCommand DataOutputBuffer dob = new DataOutputBuffer(); serializer.serialize(this, dob); return new Message(FBUtilities.getLocalAddress(), - StageManager.readStage_, + StageManager.READ_STAGE, StorageService.rangeVerbHandler_, Arrays.copyOf(dob.getData(), dob.getLength())); } diff --git a/src/java/org/apache/cassandra/db/RangeSliceCommand.java b/src/java/org/apache/cassandra/db/RangeSliceCommand.java index 0840afe0c7..6750a3ba22 100644 --- a/src/java/org/apache/cassandra/db/RangeSliceCommand.java +++ b/src/java/org/apache/cassandra/db/RangeSliceCommand.java @@ -100,7 +100,7 @@ public class RangeSliceCommand DataOutputBuffer dob = new DataOutputBuffer(); serializer.serialize(this, dob); return new Message(FBUtilities.getLocalAddress(), - StageManager.readStage_, + StageManager.READ_STAGE, StorageService.rangeSliceVerbHandler_, Arrays.copyOf(dob.getData(), dob.getLength())); } diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index 68b18b3065..25e88a40f5 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -54,7 +54,7 @@ public abstract class ReadCommand ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream(bos); ReadCommand.serializer().serialize(this, dos); - return new Message(FBUtilities.getLocalAddress(), StageManager.readStage_, StorageService.readVerbHandler_, bos.toByteArray()); + return new Message(FBUtilities.getLocalAddress(), StageManager.READ_STAGE, StorageService.readVerbHandler_, bos.toByteArray()); } public final QueryPath queryPath; diff --git a/src/java/org/apache/cassandra/db/ReadResponse.java b/src/java/org/apache/cassandra/db/ReadResponse.java index 2c2a55df77..16a9d77e5f 100644 --- a/src/java/org/apache/cassandra/db/ReadResponse.java +++ b/src/java/org/apache/cassandra/db/ReadResponse.java @@ -23,6 +23,7 @@ import java.io.DataInputStream; import java.io.DataOutputStream; import java.io.IOException; +import org.apache.cassandra.concurrent.StageManager; import org.apache.cassandra.io.ICompactSerializer; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -55,7 +56,7 @@ private static ICompactSerializer serializer_; ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream( bos ); ReadResponse.serializer().serialize(readResponse, dos); - Message message = new Message(FBUtilities.getLocalAddress(), MessagingService.responseStage_, MessagingService.responseVerbHandler_, bos.toByteArray()); + Message message = new Message(FBUtilities.getLocalAddress(), StageManager.RESPONSE_STAGE, MessagingService.responseVerbHandler_, bos.toByteArray()); return message; } diff --git a/src/java/org/apache/cassandra/db/RowMutation.java b/src/java/org/apache/cassandra/db/RowMutation.java index d67f53a301..5ddb8b1668 100644 --- a/src/java/org/apache/cassandra/db/RowMutation.java +++ b/src/java/org/apache/cassandra/db/RowMutation.java @@ -216,7 +216,7 @@ public class RowMutation ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream(bos); serializer().serialize(this, dos); - return new Message(FBUtilities.getLocalAddress(), StageManager.mutationStage_, verbHandlerName, bos.toByteArray()); + return new Message(FBUtilities.getLocalAddress(), StageManager.MUTATION_STAGE, verbHandlerName, bos.toByteArray()); } public static RowMutation getRowMutationFromMutations(String keyspace, String key, Map> cfmap) diff --git a/src/java/org/apache/cassandra/db/RowMutationMessage.java b/src/java/org/apache/cassandra/db/RowMutationMessage.java index 1b74bd838d..1e4de769b8 100644 --- a/src/java/org/apache/cassandra/db/RowMutationMessage.java +++ b/src/java/org/apache/cassandra/db/RowMutationMessage.java @@ -51,7 +51,7 @@ public class RowMutationMessage ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream( bos ); RowMutationMessage.serializer().serialize(this, dos); - return new Message(FBUtilities.getLocalAddress(), StageManager.mutationStage_, verbHandlerName, bos.toByteArray()); + return new Message(FBUtilities.getLocalAddress(), StageManager.MUTATION_STAGE, verbHandlerName, bos.toByteArray()); } @XmlElement(name="RowMutation") diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index 900b93d3ca..4b4fd5a645 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -22,7 +22,6 @@ import java.io.*; import java.util.*; import java.net.InetAddress; -import org.apache.cassandra.concurrent.SingleThreadedStage; import org.apache.cassandra.concurrent.StageManager; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.net.IVerbHandler; @@ -99,8 +98,6 @@ public class Gossiper implements IFailureDetectionEventListener, IEndPointStateC } final static int MAX_GOSSIP_PACKET_SIZE = 1428; - /* GS - abbreviation for GOSSIPER_STAGE */ - final static String GOSSIP_STAGE = "GS"; /* GSV - abbreviation for GOSSIP-DIGEST-SYN-VERB */ final static String JOIN_VERB_HANDLER = "JVH"; /* GSV - abbreviation for GOSSIP-DIGEST-SYN-VERB */ @@ -153,8 +150,6 @@ public class Gossiper implements IFailureDetectionEventListener, IEndPointStateC MessagingService.instance().registerVerbHandlers(GOSSIP_DIGEST_SYN_VERB, new GossipDigestSynVerbHandler()); MessagingService.instance().registerVerbHandlers(GOSSIP_DIGEST_ACK_VERB, new GossipDigestAckVerbHandler()); MessagingService.instance().registerVerbHandlers(GOSSIP_DIGEST_ACK2_VERB, new GossipDigestAck2VerbHandler()); - /* register the Gossip stage */ - StageManager.registerStage( Gossiper.GOSSIP_STAGE, new SingleThreadedStage("GMFD") ); } /** Register with the Gossiper for EndPointState notifications */ @@ -285,7 +280,7 @@ public class Gossiper implements IFailureDetectionEventListener, IEndPointStateC ByteArrayOutputStream bos = new ByteArrayOutputStream(Gossiper.MAX_GOSSIP_PACKET_SIZE); DataOutputStream dos = new DataOutputStream( bos ); GossipDigestSynMessage.serializer().serialize(gDigestMessage, dos); - return new Message(localEndPoint_, Gossiper.GOSSIP_STAGE, GOSSIP_DIGEST_SYN_VERB, bos.toByteArray()); + return new Message(localEndPoint_, StageManager.GOSSIP_STAGE, GOSSIP_DIGEST_SYN_VERB, bos.toByteArray()); } Message makeGossipDigestAckMessage(GossipDigestAckMessage gDigestAckMessage) throws IOException @@ -295,7 +290,7 @@ public class Gossiper implements IFailureDetectionEventListener, IEndPointStateC GossipDigestAckMessage.serializer().serialize(gDigestAckMessage, dos); if (logger_.isTraceEnabled()) logger_.trace("@@@@ Size of GossipDigestAckMessage is " + bos.toByteArray().length); - return new Message(localEndPoint_, Gossiper.GOSSIP_STAGE, GOSSIP_DIGEST_ACK_VERB, bos.toByteArray()); + return new Message(localEndPoint_, StageManager.GOSSIP_STAGE, GOSSIP_DIGEST_ACK_VERB, bos.toByteArray()); } Message makeGossipDigestAck2Message(GossipDigestAck2Message gDigestAck2Message) throws IOException @@ -303,7 +298,7 @@ public class Gossiper implements IFailureDetectionEventListener, IEndPointStateC ByteArrayOutputStream bos = new ByteArrayOutputStream(Gossiper.MAX_GOSSIP_PACKET_SIZE); DataOutputStream dos = new DataOutputStream(bos); GossipDigestAck2Message.serializer().serialize(gDigestAck2Message, dos); - return new Message(localEndPoint_, Gossiper.GOSSIP_STAGE, GOSSIP_DIGEST_ACK2_VERB, bos.toByteArray()); + return new Message(localEndPoint_, StageManager.GOSSIP_STAGE, GOSSIP_DIGEST_ACK2_VERB, bos.toByteArray()); } /** diff --git a/src/java/org/apache/cassandra/io/StreamRequestMessage.java b/src/java/org/apache/cassandra/io/StreamRequestMessage.java index a1b3d01ae7..2d27f2a2cd 100644 --- a/src/java/org/apache/cassandra/io/StreamRequestMessage.java +++ b/src/java/org/apache/cassandra/io/StreamRequestMessage.java @@ -55,7 +55,7 @@ class StreamRequestMessage { throw new IOError(e); } - return new Message(FBUtilities.getLocalAddress(), StageManager.streamStage_, StorageService.streamRequestVerbHandler_, bos.toByteArray() ); + return new Message(FBUtilities.getLocalAddress(), StageManager.STREAM_STAGE, StorageService.streamRequestVerbHandler_, bos.toByteArray() ); } protected StreamRequestMetadata[] streamRequestMetadata_ = new StreamRequestMetadata[0]; diff --git a/src/java/org/apache/cassandra/net/Message.java b/src/java/org/apache/cassandra/net/Message.java index ebc2b6448a..90989ad823 100644 --- a/src/java/org/apache/cassandra/net/Message.java +++ b/src/java/org/apache/cassandra/net/Message.java @@ -24,6 +24,7 @@ import java.io.IOException; import java.util.Map; import java.net.InetAddress; +import org.apache.cassandra.concurrent.StageManager; import org.apache.cassandra.io.ICompactSerializer; public class Message @@ -120,7 +121,7 @@ public class Message // TODO should take byte[] + length so we don't have to copy to a byte[] of exactly the right len public Message getReply(InetAddress from, byte[] args) { - Header header = new Header(getMessageId(), from, MessagingService.responseStage_, MessagingService.responseVerbHandler_); + Header header = new Header(getMessageId(), from, StageManager.RESPONSE_STAGE, MessagingService.responseVerbHandler_); return new Message(header, args); } diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index 1433c1edec..8d73b9d9f0 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -51,8 +51,6 @@ public class MessagingService implements IFailureDetectionEventListener private static byte[] protocol_ = new byte[16]; /* Verb Handler for the Response */ public static final String responseVerbHandler_ = "RESPONSE"; - /* Stage for responses. */ - public static final String responseStage_ = "RESPONSE-STAGE"; /* This records all the results mapped by message Id */ private static ICachetable callbackMap_; @@ -84,7 +82,7 @@ public class MessagingService implements IFailureDetectionEventListener private static volatile MessagingService messagingService_ = new MessagingService(); - private static final int MESSAGE_DESERIALIZE_THREADS = 4; + public static final int MESSAGE_DESERIALIZE_THREADS = 4; public static int getVersion() { @@ -121,34 +119,33 @@ public class MessagingService implements IFailureDetectionEventListener * before the callback is evicted from the table. The concurrency level is set at 128 * which is the sum of the threads in the pool that adds shit into the table and the * pool that retrives the callback from here. - */ - int maxSize = MESSAGE_DESERIALIZE_THREADS; + */ callbackMap_ = new Cachetable( 2 * DatabaseDescriptor.getRpcTimeout() ); taskCompletionMap_ = new Cachetable( 2 * DatabaseDescriptor.getRpcTimeout() ); - messageDeserializationExecutor_ = new JMXEnabledThreadPoolExecutor( maxSize, - maxSize, + messageDeserializationExecutor_ = new JMXEnabledThreadPoolExecutor( + MESSAGE_DESERIALIZE_THREADS, + MESSAGE_DESERIALIZE_THREADS, Integer.MAX_VALUE, TimeUnit.SECONDS, new LinkedBlockingQueue(), new NamedThreadFactory("MESSAGING-SERVICE-POOL") - ); + ); - messageDeserializerExecutor_ = new JMXEnabledThreadPoolExecutor( maxSize, - maxSize, + messageDeserializerExecutor_ = new JMXEnabledThreadPoolExecutor( + MESSAGE_DESERIALIZE_THREADS, + MESSAGE_DESERIALIZE_THREADS, Integer.MAX_VALUE, TimeUnit.SECONDS, new LinkedBlockingQueue(), new NamedThreadFactory("MESSAGE-DESERIALIZER-POOL") - ); + ); streamExecutor_ = new JMXEnabledThreadPoolExecutor("MESSAGE-STREAMING-POOL"); protocol_ = hash(HashingSchemes.MD5, "FB-MESSAGING".getBytes()); /* register the response verb handler */ registerVerbHandlers(MessagingService.responseVerbHandler_, new ResponseVerbHandler()); - /* register stage for response */ - StageManager.registerStage(MessagingService.responseStage_, new MultiThreadedStage("RESPONSE-STAGE", maxSize) ); } public byte[] hash(String type, byte data[]) diff --git a/src/java/org/apache/cassandra/net/TcpConnectionManager.java b/src/java/org/apache/cassandra/net/TcpConnectionManager.java index bcd8383105..8cdfef2866 100644 --- a/src/java/org/apache/cassandra/net/TcpConnectionManager.java +++ b/src/java/org/apache/cassandra/net/TcpConnectionManager.java @@ -19,11 +19,9 @@ package org.apache.cassandra.net; import java.io.IOException; -import java.util.*; -import java.util.concurrent.locks.*; import java.net.InetAddress; -import org.apache.log4j.Logger; +import org.apache.cassandra.concurrent.StageManager; class TcpConnectionManager { @@ -49,7 +47,7 @@ class TcpConnectionManager */ synchronized TcpConnection getConnection(Message msg) throws IOException { - if (MessagingService.responseStage_.equals(msg.getMessageType())) + if (StageManager.RESPONSE_STAGE.equals(msg.getMessageType())) { if (ackCon == null) ackCon = newCon(); diff --git a/src/java/org/apache/cassandra/service/AntiEntropyService.java b/src/java/org/apache/cassandra/service/AntiEntropyService.java index d866fd8864..cfa6d459f5 100644 --- a/src/java/org/apache/cassandra/service/AntiEntropyService.java +++ b/src/java/org/apache/cassandra/service/AntiEntropyService.java @@ -23,7 +23,6 @@ import java.net.InetAddress; import java.util.*; import java.util.concurrent.*; -import org.apache.cassandra.concurrent.SingleThreadedStage; import org.apache.cassandra.concurrent.StageManager; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.CompactionManager; @@ -90,7 +89,6 @@ public class AntiEntropyService { private static final Logger logger = Logger.getLogger(AntiEntropyService.class); - public final static String AE_SERVICE_STAGE = "AE-SERVICE-STAGE"; public final static String TREE_REQUEST_VERB = "TREE-REQUEST-VERB"; public final static String TREE_RESPONSE_VERB = "TREE-RESPONSE-VERB"; @@ -137,8 +135,6 @@ public class AntiEntropyService */ private AntiEntropyService() { - StageManager.registerStage(AE_SERVICE_STAGE, new SingleThreadedStage(AE_SERVICE_STAGE)); - MessagingService.instance().registerVerbHandlers(TREE_REQUEST_VERB, new TreeRequestVerbHandler()); MessagingService.instance().registerVerbHandlers(TREE_RESPONSE_VERB, new TreeResponseVerbHandler()); naturalRepairs = new ConcurrentHashMap(); @@ -230,7 +226,7 @@ public class AntiEntropyService for (Differencer differencer : differencers) { logger.info("Queueing comparison " + differencer); - StageManager.getStage(AE_SERVICE_STAGE).execute(differencer); + StageManager.getStage(StageManager.AE_SERVICE_STAGE).execute(differencer); } } @@ -487,7 +483,7 @@ public class AntiEntropyService for (MerkleTree.RowHash minrow : minrows) range.addHash(minrow); - StageManager.getStage(AE_SERVICE_STAGE).execute(this); + StageManager.getStage(StageManager.AE_SERVICE_STAGE).execute(this); logger.debug("Validated " + validated + " rows into AEService tree for " + cf); } @@ -681,7 +677,7 @@ public class AntiEntropyService ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream(bos); SERIALIZER.serialize(new CFPair(table, cf), dos); - return new Message(FBUtilities.getLocalAddress(), AE_SERVICE_STAGE, TREE_REQUEST_VERB, bos.toByteArray()); + return new Message(FBUtilities.getLocalAddress(), StageManager.AE_SERVICE_STAGE, TREE_REQUEST_VERB, bos.toByteArray()); } catch(IOException e) { @@ -739,7 +735,7 @@ public class AntiEntropyService ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream(bos); SERIALIZER.serialize(validator, dos); - return new Message(local, AE_SERVICE_STAGE, TREE_RESPONSE_VERB, bos.toByteArray()); + return new Message(local, StageManager.AE_SERVICE_STAGE, TREE_RESPONSE_VERB, bos.toByteArray()); } catch(IOException e) { diff --git a/src/java/org/apache/cassandra/service/StorageLoadBalancer.java b/src/java/org/apache/cassandra/service/StorageLoadBalancer.java index 01f48fe8f2..1c70682849 100644 --- a/src/java/org/apache/cassandra/service/StorageLoadBalancer.java +++ b/src/java/org/apache/cassandra/service/StorageLoadBalancer.java @@ -177,7 +177,6 @@ public final class StorageLoadBalancer implements IEndPointStateChangeSubscriber private static final Logger logger_ = Logger.getLogger(StorageLoadBalancer.class); - private static final String lbStage_ = "LOAD-BALANCER-STAGE"; private static final String moveMessageVerbHandler_ = "MOVE-MESSAGE-VERB-HANDLER"; /* time to delay in minutes the actual load balance procedure if heavily loaded */ private static final int delay_ = 5; @@ -199,7 +198,6 @@ public final class StorageLoadBalancer implements IEndPointStateChangeSubscriber private StorageLoadBalancer() { - StageManager.registerStage(StorageLoadBalancer.lbStage_, new SingleThreadedStage(StorageLoadBalancer.lbStage_)); MessagingService.instance().registerVerbHandlers(StorageLoadBalancer.moveMessageVerbHandler_, new MoveMessageVerbHandler()); Gossiper.instance().register(this); } diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 5dd38a0652..0c3bf5e142 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -19,7 +19,6 @@ package org.apache.cassandra.service; import java.io.ByteArrayInputStream; import java.io.DataInputStream; -import java.io.IOError; import java.io.IOException; import java.util.*; import java.util.concurrent.TimeUnit; @@ -130,7 +129,7 @@ public class StorageProxy implements StorageProxyMBean rm.apply(); } }; - StageManager.getStage(StageManager.mutationStage_).execute(runnable); + StageManager.getStage(StageManager.MUTATION_STAGE).execute(runnable); } else { @@ -270,7 +269,7 @@ public class StorageProxy implements StorageProxyMBean responseHandler.localResponse(); } }; - StageManager.getStage(StageManager.mutationStage_).execute(runnable); + StageManager.getStage(StageManager.MUTATION_STAGE).execute(runnable); } private static int determineBlockFor(int naturalTargets, int hintedTargets, ConsistencyLevel consistency_level) @@ -497,7 +496,7 @@ public class StorageProxy implements StorageProxyMBean for (ReadCommand command: commands) { Callable callable = new weakReadLocalCallable(command); - futures.add(StageManager.getStage(StageManager.readStage_).execute(callable)); + futures.add(StageManager.getStage(StageManager.READ_STAGE).execute(callable)); } for (Future future : futures) { diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 6f6fe8535c..0272ce874f 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -1329,7 +1329,7 @@ public final class StorageService implements IEndPointStateChangeSubscriber, Sto } } }; - StageManager.getStage(StageManager.streamStage_).execute(new Runnable() + StageManager.getStage(StageManager.STREAM_STAGE).execute(new Runnable() { public void run() { diff --git a/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java b/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java index e48968a6a5..2abe571ab3 100644 --- a/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java +++ b/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java @@ -235,7 +235,7 @@ public class AntiEntropyServiceTest extends CleanupHelper Future flushAES() { - return StageManager.getStage(AE_SERVICE_STAGE).execute(new Callable(){ + return StageManager.getStage(StageManager.AE_SERVICE_STAGE).execute(new Callable(){ public Boolean call() { return true;