[infinispan-commits] Infinispan SVN: r1669 - in trunk/server: core/src/main/scala/org/infinispan/server/core/transport/netty and 2 other directories.
infinispan-commits at lists.jboss.org
infinispan-commits at lists.jboss.org
Wed Apr 7 11:06:43 EDT 2010
Author: galder.zamarreno at jboss.com
Date: 2010-04-07 11:06:42 -0400 (Wed, 07 Apr 2010)
New Revision: 1669
Modified:
trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala
trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala
trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala
trunk/server/memcached/src/main/scala/org/infinispan/server/memcached/MemcachedServer.scala
Log:
[ISPN-394] (Add proper thread naming for server worker threads) Done.
Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala 2010-04-07 13:44:35 UTC (rev 1668)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/AbstractProtocolServer.scala 2010-04-07 15:06:42 UTC (rev 1669)
@@ -10,7 +10,7 @@
* @author Galder Zamarreño
* @since 4.1
*/
-abstract class AbstractProtocolServer extends ProtocolServer {
+abstract class AbstractProtocolServer(threadNamePrefix: String) extends ProtocolServer {
private var host: String = _
private var port: Int = _
private var masterThreads: Int = _
@@ -32,8 +32,7 @@
// TODO: ... requests such as when the lenght of data is bigger than the expected data itself.
val nettyEncoder = if (encoder != null) new EncoderAdapter(encoder) else null
val address = new InetSocketAddress(host, port)
- // TODO change cache name 'default' to something more meaningful and dependent of protocol
- transport = new NettyTransport(this, nettyEncoder, address, masterThreads, workerThreads, "default")
+ transport = new NettyTransport(this, nettyEncoder, address, masterThreads, workerThreads, threadNamePrefix)
transport.start
}
Modified: trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala
===================================================================
--- trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala 2010-04-07 13:44:35 UTC (rev 1668)
+++ trunk/server/core/src/main/scala/org/infinispan/server/core/transport/netty/NettyTransport.scala 2010-04-07 15:06:42 UTC (rev 1669)
@@ -11,6 +11,7 @@
import scala.collection.JavaConversions._
import org.infinispan.manager.CacheManager
import org.infinispan.server.core.{ProtocolServer, Logging}
+import org.jboss.netty.util.{ThreadNameDeterminer, ThreadRenamingRunnable}
/**
* // TODO: Document this
@@ -19,11 +20,11 @@
*/
class NettyTransport(server: ProtocolServer, encoder: ChannelDownstreamHandler,
address: SocketAddress, masterThreads: Int, workerThreads: Int,
- cacheName: String) extends Transport {
+ threadNamePrefix: String) extends Transport {
import NettyTransport._
- val serverChannels = new DefaultChannelGroup(cacheName + "-channels")
- val acceptedChannels = new DefaultChannelGroup(cacheName + "-accepted")
+ val serverChannels = new DefaultChannelGroup(threadNamePrefix + "-Channels")
+ val acceptedChannels = new DefaultChannelGroup(threadNamePrefix + "-Accepted")
val pipeline = new NettyChannelPipelineFactory(server, encoder)
val factory = {
if (workerThreads == 0)
@@ -33,29 +34,34 @@
}
lazy val masterExecutor = {
- val tf = new NamedThreadFactory(cacheName + '-' + "Master")
if (masterThreads == 0) {
debug("Configured unlimited threads for master thread pool")
- Executors.newCachedThreadPool(tf)
+ Executors.newCachedThreadPool
} else {
debug("Configured {0} threads for master thread pool", masterThreads)
- Executors.newFixedThreadPool(masterThreads, tf)
+ Executors.newFixedThreadPool(masterThreads)
}
}
lazy val workerExecutor = {
- val tf = new NamedThreadFactory(cacheName + '-' + "Worker")
if (workerThreads == 0) {
debug("Configured unlimited threads for worker thread pool")
- Executors.newCachedThreadPool(tf)
+ Executors.newCachedThreadPool
}
else {
debug("Configured {0} threads for worker thread pool", workerThreads)
- Executors.newFixedThreadPool(workerThreads, tf)
+ Executors.newFixedThreadPool(masterThreads)
}
}
override def start {
+ ThreadRenamingRunnable.setThreadNameDeterminer(new ThreadNameDeterminer {
+ override def determineThreadName(currentThreadName: String, proposedThreadName: String): String = {
+ val index = proposedThreadName.findIndexOf(_ == '#')
+ val typeInFix = if (proposedThreadName.contains("boss")) "Master-" else "Worker-"
+ threadNamePrefix + typeInFix + proposedThreadName.substring(index + 1, proposedThreadName.length)
+ }
+ })
val bootstrap = new ServerBootstrap(factory);
bootstrap.setPipelineFactory(pipeline);
val ch = bootstrap.bind(address);
Modified: trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala
===================================================================
--- trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala 2010-04-07 13:44:35 UTC (rev 1668)
+++ trunk/server/hotrod/src/main/scala/org/infinispan/server/hotrod/HotRodServer.scala 2010-04-07 15:06:42 UTC (rev 1669)
@@ -10,7 +10,7 @@
* @since 4.1
*/
-class HotRodServer extends AbstractProtocolServer {
+class HotRodServer extends AbstractProtocolServer("HotRod") {
override def getEncoder: Encoder = new HotRodEncoder
Modified: trunk/server/memcached/src/main/scala/org/infinispan/server/memcached/MemcachedServer.scala
===================================================================
--- trunk/server/memcached/src/main/scala/org/infinispan/server/memcached/MemcachedServer.scala 2010-04-07 13:44:35 UTC (rev 1668)
+++ trunk/server/memcached/src/main/scala/org/infinispan/server/memcached/MemcachedServer.scala 2010-04-07 15:06:42 UTC (rev 1669)
@@ -11,7 +11,7 @@
* @since
*/
-class MemcachedServer extends AbstractProtocolServer {
+class MemcachedServer extends AbstractProtocolServer("Memcached") {
protected lazy val scheduler = Executors.newScheduledThreadPool(1)
More information about the infinispan-commits
mailing list