diff --git a/CHANGES.txt b/CHANGES.txt
index 6d77f3c78f..4954b7f7b5 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -5,6 +5,7 @@
* Always reject inequality on the partition key without token()
(CASSANDRA-7722)
* Always send Paxos commit to all replicas (CASSANDRA-7479)
+ * Make disruptor_thrift_server invocation pool configurable (CASSANDRA-7594)
2.0.10
diff --git a/build.xml b/build.xml
index dd59bd28aa..f456fa88d2 100644
--- a/build.xml
+++ b/build.xml
@@ -361,7 +361,7 @@
-
+
@@ -467,7 +467,7 @@
-
+
diff --git a/lib/thrift-server-internal-only-0.3.3.jar b/lib/thrift-server-0.3.6.jar
similarity index 50%
rename from lib/thrift-server-internal-only-0.3.3.jar
rename to lib/thrift-server-0.3.6.jar
index 6a1fbae73d..c974f75e4e 100644
Binary files a/lib/thrift-server-internal-only-0.3.3.jar and b/lib/thrift-server-0.3.6.jar differ
diff --git a/src/java/org/apache/cassandra/thrift/THsHaDisruptorServer.java b/src/java/org/apache/cassandra/thrift/THsHaDisruptorServer.java
index e3b89d26cf..dd501ec5e7 100644
--- a/src/java/org/apache/cassandra/thrift/THsHaDisruptorServer.java
+++ b/src/java/org/apache/cassandra/thrift/THsHaDisruptorServer.java
@@ -19,9 +19,14 @@
package org.apache.cassandra.thrift;
import java.net.InetSocketAddress;
+import java.util.concurrent.SynchronousQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
import com.thinkaurelius.thrift.Message;
import com.thinkaurelius.thrift.TDisruptorServer;
+import org.apache.cassandra.concurrent.JMXEnabledThreadPoolExecutor;
+import org.apache.cassandra.concurrent.NamedThreadFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -78,6 +83,13 @@ public class THsHaDisruptorServer extends TDisruptorServer
throw new RuntimeException(String.format("Unable to create thrift socket to %s:%s", addr.getAddress(), addr.getPort()), e);
}
+ ThreadPoolExecutor invoker = new JMXEnabledThreadPoolExecutor(DatabaseDescriptor.getRpcMinThreads(),
+ DatabaseDescriptor.getRpcMaxThreads(),
+ 60L,
+ TimeUnit.SECONDS,
+ new SynchronousQueue(),
+ new NamedThreadFactory("RPC-Thread"), "RPC-THREAD-POOL");
+
com.thinkaurelius.thrift.util.TBinaryProtocol.Factory protocolFactory = new com.thinkaurelius.thrift.util.TBinaryProtocol.Factory(true, true);
TDisruptorServer.Args serverArgs = new TDisruptorServer.Args(serverTransport).useHeapBasedAllocation(true)
@@ -87,6 +99,7 @@ public class THsHaDisruptorServer extends TDisruptorServer
.outputProtocolFactory(protocolFactory)
.processor(args.processor)
.maxFrameSizeInBytes(DatabaseDescriptor.getThriftFramedTransportSize())
+ .invocationExecutor(invoker)
.alwaysReallocateBuffers(true);
return new THsHaDisruptorServer(serverArgs);