From 3ed4e051b89b7fa6e42d665cf1f47cdda13175bc Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Wed, 6 May 2009 02:27:54 +0000 Subject: [PATCH] add pending tasks mbean to all DebuggableTPE. clean out ad-hoc queue length logging. patch by jbellis; reviewed by Eric Evans for CASSANDRA-135 git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@772030 13f79535-47bb-0310-9956-ffa450edef68 --- .../concurrent/ContinuationStage.java | 2 +- .../DebuggableThreadPoolExecutor.java | 26 +++++++-- .../DebuggableThreadPoolExecutorMBean.java | 6 ++ .../apache/cassandra/concurrent/IStage.java | 2 +- .../concurrent/MultiThreadedStage.java | 4 +- .../SingleThreadedContinuationStage.java | 2 +- .../concurrent/SingleThreadedStage.java | 5 +- .../cassandra/concurrent/StageManager.java | 4 +- .../cassandra/db/ColumnFamilyStore.java | 5 -- .../cassandra/db/ColumnFamilyStoreMBean.java | 5 -- .../org/apache/cassandra/db/Memtable.java | 16 ++---- .../net/MessageDeserializationTask.java | 27 ++++----- .../net/MessageSerializationTask.java | 5 -- .../cassandra/net/MessagingService.java | 55 ++++++------------- .../cassandra/net/MessagingServiceMBean.java | 29 ---------- .../apache/cassandra/net/TcpConnection.java | 3 +- .../apache/cassandra/net/UdpConnection.java | 2 +- .../cassandra/service/StorageService.java | 1 - 18 files changed, 71 insertions(+), 128 deletions(-) create mode 100644 src/java/org/apache/cassandra/concurrent/DebuggableThreadPoolExecutorMBean.java delete mode 100644 src/java/org/apache/cassandra/net/MessagingServiceMBean.java diff --git a/src/java/org/apache/cassandra/concurrent/ContinuationStage.java b/src/java/org/apache/cassandra/concurrent/ContinuationStage.java index 0c5291d646..d8bdfc037f 100644 --- a/src/java/org/apache/cassandra/concurrent/ContinuationStage.java +++ b/src/java/org/apache/cassandra/concurrent/ContinuationStage.java @@ -85,7 +85,7 @@ public class ContinuationStage implements IStage return executorService_.isShutdown(); } - public long getTaskCount(){ + public long getPendingTasks(){ return (executorService_.getTaskCount() - executorService_.getCompletedTaskCount()); } } diff --git a/src/java/org/apache/cassandra/concurrent/DebuggableThreadPoolExecutor.java b/src/java/org/apache/cassandra/concurrent/DebuggableThreadPoolExecutor.java index b72ed35696..51faa4c05f 100644 --- a/src/java/org/apache/cassandra/concurrent/DebuggableThreadPoolExecutor.java +++ b/src/java/org/apache/cassandra/concurrent/DebuggableThreadPoolExecutor.java @@ -19,10 +19,12 @@ package org.apache.cassandra.concurrent; import java.util.concurrent.*; +import java.lang.management.ManagementFactory; -import org.apache.cassandra.utils.LogUtil; import org.apache.log4j.Logger; -import org.apache.cassandra.utils.*; + +import javax.management.MBeanServer; +import javax.management.ObjectName; /** * This is a wrapper class for the ScheduledThreadPoolExecutor. It provides an implementation @@ -32,7 +34,7 @@ import org.apache.cassandra.utils.*; * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) */ -public class DebuggableThreadPoolExecutor extends ThreadPoolExecutor +public class DebuggableThreadPoolExecutor extends ThreadPoolExecutor implements DebuggableThreadPoolExecutorMBean { private static Logger logger_ = Logger.getLogger(DebuggableThreadPoolExecutor.class); @@ -46,12 +48,26 @@ public class DebuggableThreadPoolExecutor extends ThreadPoolExecutor long keepAliveTime, TimeUnit unit, BlockingQueue workQueue, - ThreadFactory threadFactory) + ThreadFactoryImpl threadFactory) { super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory); super.prestartAllCoreThreads(); + MBeanServer mbs = ManagementFactory.getPlatformMBeanServer(); + try + { + mbs.registerMBean(this, new ObjectName("org.apache.cassandra.concurrent:type=" + threadFactory.id_)); + } + catch (Exception e) + { + throw new RuntimeException(e); + } } - + + public long getPendingTasks() + { + return getTaskCount() - getCompletedTaskCount(); + } + /* * * (non-Javadoc) diff --git a/src/java/org/apache/cassandra/concurrent/DebuggableThreadPoolExecutorMBean.java b/src/java/org/apache/cassandra/concurrent/DebuggableThreadPoolExecutorMBean.java new file mode 100644 index 0000000000..798a888b51 --- /dev/null +++ b/src/java/org/apache/cassandra/concurrent/DebuggableThreadPoolExecutorMBean.java @@ -0,0 +1,6 @@ +package org.apache.cassandra.concurrent; + +public interface DebuggableThreadPoolExecutorMBean +{ + public long getPendingTasks(); +} diff --git a/src/java/org/apache/cassandra/concurrent/IStage.java b/src/java/org/apache/cassandra/concurrent/IStage.java index 1ed9fc0625..f33e3cb228 100644 --- a/src/java/org/apache/cassandra/concurrent/IStage.java +++ b/src/java/org/apache/cassandra/concurrent/IStage.java @@ -116,5 +116,5 @@ public interface IStage * pending on this stage to be executed. * @return task count. */ - public long getTaskCount(); + public long getPendingTasks(); } diff --git a/src/java/org/apache/cassandra/concurrent/MultiThreadedStage.java b/src/java/org/apache/cassandra/concurrent/MultiThreadedStage.java index 75a27896b4..bbdcad62c3 100644 --- a/src/java/org/apache/cassandra/concurrent/MultiThreadedStage.java +++ b/src/java/org/apache/cassandra/concurrent/MultiThreadedStage.java @@ -92,7 +92,7 @@ public class MultiThreadedStage implements IStage return executorService_.isShutdown(); } - public long getTaskCount(){ - return (executorService_.getTaskCount() - executorService_.getCompletedTaskCount()); + public long getPendingTasks(){ + return executorService_.getPendingTasks(); } } diff --git a/src/java/org/apache/cassandra/concurrent/SingleThreadedContinuationStage.java b/src/java/org/apache/cassandra/concurrent/SingleThreadedContinuationStage.java index 7a4cd33877..fa1fcff2fe 100644 --- a/src/java/org/apache/cassandra/concurrent/SingleThreadedContinuationStage.java +++ b/src/java/org/apache/cassandra/concurrent/SingleThreadedContinuationStage.java @@ -92,7 +92,7 @@ public class SingleThreadedContinuationStage implements IStage return executorService_.isShutdown(); } - public long getTaskCount(){ + public long getPendingTasks(){ return (executorService_.getTaskCount() - executorService_.getCompletedTaskCount()); } /* Finished implementing the IStage interface methods */ diff --git a/src/java/org/apache/cassandra/concurrent/SingleThreadedStage.java b/src/java/org/apache/cassandra/concurrent/SingleThreadedStage.java index 27ae871955..24480d33f6 100644 --- a/src/java/org/apache/cassandra/concurrent/SingleThreadedStage.java +++ b/src/java/org/apache/cassandra/concurrent/SingleThreadedStage.java @@ -94,8 +94,7 @@ public class SingleThreadedStage implements IStage return executorService_.isShutdown(); } - public long getTaskCount(){ - return (executorService_.getTaskCount() - executorService_.getCompletedTaskCount()); + public long getPendingTasks(){ + return executorService_.getPendingTasks(); } - /* Finished implementing the IStage interface methods */ } diff --git a/src/java/org/apache/cassandra/concurrent/StageManager.java b/src/java/org/apache/cassandra/concurrent/StageManager.java index c3fc2840d1..f771588d6d 100644 --- a/src/java/org/apache/cassandra/concurrent/StageManager.java +++ b/src/java/org/apache/cassandra/concurrent/StageManager.java @@ -23,8 +23,6 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.ExecutorService; -import org.apache.cassandra.continuations.Suspendable; - /** * This class manages all stages that exist within a process. The application registers @@ -101,7 +99,7 @@ public class StageManager */ public static long getStageTaskCount(String stage) { - return stageQueues_.get(stage).getTaskCount(); + return stageQueues_.get(stage).getPendingTasks(); } /** diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 1f283b600f..2787b6316a 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -1432,11 +1432,6 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return memtableSwitchCount; } - public int getMemtableTasks() - { - return memtable_.get().getPendingTasks(); - } - /** * clears out all data associated with this ColumnFamily. * For use in testing. diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStoreMBean.java b/src/java/org/apache/cassandra/db/ColumnFamilyStoreMBean.java index 75840b849e..94e91c343d 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStoreMBean.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStoreMBean.java @@ -49,11 +49,6 @@ public interface ColumnFamilyStoreMBean */ public int getMemtableSwitchCount(); - /** - * @return the number of tasks waiting to run on the memtable executor - */ - public int getMemtableTasks(); - /** * Triggers an immediate memtable flush. */ diff --git a/src/java/org/apache/cassandra/db/Memtable.java b/src/java/org/apache/cassandra/db/Memtable.java index ae926d8fb1..9eeec349fa 100644 --- a/src/java/org/apache/cassandra/db/Memtable.java +++ b/src/java/org/apache/cassandra/db/Memtable.java @@ -45,7 +45,7 @@ public class Memtable implements Comparable { private static Logger logger_ = Logger.getLogger( Memtable.class ); private static Set runningExecutorServices_ = new NonBlockingHashSet(); - private boolean isFrozen_; + private static AtomicInteger executorCount_ = new AtomicInteger(0); public static void shutdown() { @@ -56,6 +56,7 @@ public class Memtable implements Comparable } private MemtableThreadPoolExecutor executor_; + private boolean isFrozen_; private int threshold_ = DatabaseDescriptor.getMemtableSize()*1024*1024; private int thresholdCount_ = (int)(DatabaseDescriptor.getMemtableObjectCount()*1024*1024); @@ -72,12 +73,12 @@ public class Memtable implements Comparable Memtable(String table, String cfName) { - executor_ = new MemtableThreadPoolExecutor(); - runningExecutorServices_.add(executor_); - table_ = table; cfName_ = cfName; creationTime_ = System.currentTimeMillis(); + + executor_ = new MemtableThreadPoolExecutor(); + runningExecutorServices_.add(executor_); } class Putter implements Runnable @@ -169,11 +170,6 @@ public class Memtable implements Comparable return cfName_; } - int getPendingTasks() - { - return (int)(executor_.getTaskCount() - executor_.getCompletedTaskCount()); - } - private synchronized void enqueueFlush(CommitLog.CommitLogContext cLogCtx) { if (!isFrozen_) @@ -367,7 +363,7 @@ public class Memtable implements Comparable public MemtableThreadPoolExecutor() { - super("FAST-MEMTABLE-POOL"); + super("MEMTABLE-POOL-" + cfName_ + executorCount_.addAndGet(1)); } protected void terminated() diff --git a/src/java/org/apache/cassandra/net/MessageDeserializationTask.java b/src/java/org/apache/cassandra/net/MessageDeserializationTask.java index d21e4c6b02..0cd6578151 100644 --- a/src/java/org/apache/cassandra/net/MessageDeserializationTask.java +++ b/src/java/org/apache/cassandra/net/MessageDeserializationTask.java @@ -48,23 +48,20 @@ class MessageDeserializationTask implements Runnable public void run() { - /* For DEBUG only. Printing queue length */ - DebuggableThreadPoolExecutor es = (DebuggableThreadPoolExecutor)MessagingService.getDeserilizationExecutor(); - logger_.debug( "Message Deserialization Task: " + (es.getTaskCount() - es.getCompletedTaskCount()) ); - /* END DEBUG */ + Message message = null; try - { - Message message = (Message)serializer_.deserialize(bytes_); - - if ( message != null ) - { - message = SinkManager.processServerMessageSink(message); - MessagingService.receive(message); - } + { + message = serializer_.deserialize(bytes_); } - catch ( IOException ex ) - { - logger_.warn(LogUtil.throwableToString(ex)); + catch (IOException e) + { + throw new RuntimeException(e); + } + + if ( message != null ) + { + message = SinkManager.processServerMessageSink(message); + MessagingService.receive(message); } } diff --git a/src/java/org/apache/cassandra/net/MessageSerializationTask.java b/src/java/org/apache/cassandra/net/MessageSerializationTask.java index 9cee3143f2..3078a9a356 100644 --- a/src/java/org/apache/cassandra/net/MessageSerializationTask.java +++ b/src/java/org/apache/cassandra/net/MessageSerializationTask.java @@ -52,11 +52,6 @@ class MessageSerializationTask implements Runnable public void run() { - /* For DEBUG only. Printing queue length */ - DebuggableThreadPoolExecutor es = (DebuggableThreadPoolExecutor)MessagingService.getWriteExecutor(); - logger_.debug( "Message Serialization Task: " + (es.getTaskCount() - es.getCompletedTaskCount()) ); - /* END DEBUG */ - /* Adding the message to be serialized in the TLS. For accessing in the afterExecute() */ Context ctx = new Context(); ctx.put(this.getClass().getName(), message_); diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index 33b9e04899..bdc46d1f21 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -18,43 +18,32 @@ package org.apache.cassandra.net; -import java.io.*; -import java.lang.management.ManagementFactory; -import java.net.*; -import java.security.MessageDigest; -import java.util.*; -import java.nio.ByteBuffer; -import java.util.concurrent.*; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.locks.ReentrantLock; -import java.nio.channels.*; import org.apache.cassandra.concurrent.*; -import org.apache.cassandra.net.io.*; -import org.apache.cassandra.utils.*; -import javax.management.MBeanServer; -import javax.management.ObjectName; -import javax.xml.bind.*; -import org.apache.cassandra.concurrent.DebuggableThreadPoolExecutor; -import org.apache.cassandra.concurrent.IStage; -import org.apache.cassandra.concurrent.MultiThreadedStage; -import org.apache.cassandra.concurrent.StageManager; -import org.apache.cassandra.concurrent.ThreadFactoryImpl; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.net.http.HttpConnectionHandler; import org.apache.cassandra.net.io.SerializerType; -import org.apache.cassandra.net.sink.SinkManager; -import org.apache.cassandra.utils.Cachetable; -import org.apache.cassandra.utils.GuidGenerator; -import org.apache.cassandra.utils.HashingSchemes; -import org.apache.cassandra.utils.ICachetable; -import org.apache.cassandra.utils.LogUtil; +import org.apache.cassandra.utils.*; import org.apache.log4j.Logger; +import java.io.IOException; +import java.net.MulticastSocket; +import java.net.ServerSocket; +import java.nio.ByteBuffer; +import java.nio.channels.SelectionKey; +import java.nio.channels.ServerSocketChannel; +import java.security.MessageDigest; +import java.util.*; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.locks.ReentrantLock; + /** * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) */ -public class MessagingService implements IMessagingService, MessagingServiceMBean +public class MessagingService implements IMessagingService { private static boolean debugOn_ = false; @@ -243,18 +232,6 @@ public class MessagingService implements IMessagingService, MessagingServiceMBea return result; } - public long getMessagingSerializerTaskCount() - { - DebuggableThreadPoolExecutor dstp = (DebuggableThreadPoolExecutor)messageSerializerExecutor_; - return dstp.getTaskCount() - dstp.getCompletedTaskCount(); - } - - public long getMessagingReceiverTaskCount() - { - DebuggableThreadPoolExecutor dstp = (DebuggableThreadPoolExecutor)messageDeserializationExecutor_; - return dstp.getTaskCount() - dstp.getCompletedTaskCount(); - } - public void listen(EndPoint localEp, boolean isHttp) throws IOException { ServerSocketChannel serverChannel = ServerSocketChannel.open(); diff --git a/src/java/org/apache/cassandra/net/MessagingServiceMBean.java b/src/java/org/apache/cassandra/net/MessagingServiceMBean.java deleted file mode 100644 index 3e0a57c155..0000000000 --- a/src/java/org/apache/cassandra/net/MessagingServiceMBean.java +++ /dev/null @@ -1,29 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.cassandra.net; - -/** - * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) - */ - -public interface MessagingServiceMBean -{ - public long getMessagingSerializerTaskCount(); - public long getMessagingReceiverTaskCount(); -} diff --git a/src/java/org/apache/cassandra/net/TcpConnection.java b/src/java/org/apache/cassandra/net/TcpConnection.java index ebd507d5f7..814039fcf9 100644 --- a/src/java/org/apache/cassandra/net/TcpConnection.java +++ b/src/java/org/apache/cassandra/net/TcpConnection.java @@ -174,7 +174,7 @@ public class TcpConnection extends SelectionKeyHandler implements Comparable byte[] data = serializer_.serialize(message); if ( data.length > 0 ) { - boolean listening = ( message.getFrom().equals(EndPoint.randomLocalEndPoint_) ) ? false : true; + boolean listening = !message.getFrom().equals(EndPoint.randomLocalEndPoint_); ByteBuffer buffer = MessagingService.packIt( data , false, false, listening); synchronized(this) { @@ -184,7 +184,6 @@ public class TcpConnection extends SelectionKeyHandler implements Comparable return; } - logger_.debug("Sending packets of size " + data.length); socketChannel_.write(buffer); if (buffer.remaining() > 0) diff --git a/src/java/org/apache/cassandra/net/UdpConnection.java b/src/java/org/apache/cassandra/net/UdpConnection.java index 0ea2aecc05..53fddbb88b 100644 --- a/src/java/org/apache/cassandra/net/UdpConnection.java +++ b/src/java/org/apache/cassandra/net/UdpConnection.java @@ -78,7 +78,7 @@ public class UdpConnection extends SelectionKeyHandler byte[] data = bos.toByteArray(); if ( data.length > 0 ) { - logger_.debug("Size of Gossip packet " + data.length); + logger_.trace("Size of Gossip packet " + data.length); byte[] protocol = BasicUtilities.intToByteArray(protocol_); ByteBuffer buffer = ByteBuffer.allocate(data.length + protocol.length); buffer.put( protocol ); diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index f8b71e692d..8840685441 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -263,7 +263,6 @@ public final class StorageService implements IEndPointStateChangeSubscriber, Sto */ private void init() { - // Register this instance with JMX try { MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();