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