[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