diff --git a/CHANGES.txt b/CHANGES.txt index 75f894a70f..245f3240f5 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -83,6 +83,7 @@ * Log failed host when preparing incremental repair (CASSANDRA-8228) * Force config client mode in CQLSSTableWriter (CASSANDRA-8281) Merged from 2.0: + * Move all hints related tasks to hints internal executor (CASSANDRA-8285) * Fix paging for multi-partition IN queries (CASSANDRA-8408) * Fix MOVED_NODE topology event never being emitted when a node moves its token (CASSANDRA-8373) diff --git a/src/java/org/apache/cassandra/concurrent/DebuggableScheduledThreadPoolExecutor.java b/src/java/org/apache/cassandra/concurrent/DebuggableScheduledThreadPoolExecutor.java index 4fc1d6c990..a301923503 100644 --- a/src/java/org/apache/cassandra/concurrent/DebuggableScheduledThreadPoolExecutor.java +++ b/src/java/org/apache/cassandra/concurrent/DebuggableScheduledThreadPoolExecutor.java @@ -35,6 +35,11 @@ public class DebuggableScheduledThreadPoolExecutor extends ScheduledThreadPoolEx super(corePoolSize, new NamedThreadFactory(threadPoolName, priority)); } + public DebuggableScheduledThreadPoolExecutor(int corePoolSize, ThreadFactory threadFactory) + { + super(corePoolSize, threadFactory); + } + public DebuggableScheduledThreadPoolExecutor(String threadPoolName) { this(1, threadPoolName, Thread.NORM_PRIORITY); diff --git a/src/java/org/apache/cassandra/concurrent/JMXEnabledScheduledThreadPoolExecutor.java b/src/java/org/apache/cassandra/concurrent/JMXEnabledScheduledThreadPoolExecutor.java new file mode 100644 index 0000000000..64d9267770 --- /dev/null +++ b/src/java/org/apache/cassandra/concurrent/JMXEnabledScheduledThreadPoolExecutor.java @@ -0,0 +1,137 @@ +/* + * 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.concurrent; + +import java.lang.management.ManagementFactory; +import java.util.List; + +import javax.management.MBeanServer; +import javax.management.ObjectName; + +import org.apache.cassandra.metrics.ThreadPoolMetrics; + +/** + * A JMX enabled wrapper for DebuggableScheduledThreadPoolExecutor. + */ +public class JMXEnabledScheduledThreadPoolExecutor extends DebuggableScheduledThreadPoolExecutor implements JMXEnabledScheduledThreadPoolExecutorMBean +{ + private final String mbeanName; + private final ThreadPoolMetrics metrics; + + public JMXEnabledScheduledThreadPoolExecutor(int corePoolSize, NamedThreadFactory threadFactory, String jmxPath) + { + super(corePoolSize, threadFactory); + + metrics = new ThreadPoolMetrics(this, jmxPath, threadFactory.id); + + MBeanServer mbs = ManagementFactory.getPlatformMBeanServer(); + mbeanName = "org.apache.cassandra." + jmxPath + ":type=" + threadFactory.id; + + try + { + mbs.registerMBean(this, new ObjectName(mbeanName)); + } + catch (Exception e) + { + throw new RuntimeException(e); + } + } + + private void unregisterMBean() + { + try + { + ManagementFactory.getPlatformMBeanServer().unregisterMBean(new ObjectName(mbeanName)); + } + catch (Exception e) + { + throw new RuntimeException(e); + } + + // release metrics + metrics.release(); + } + + @Override + public synchronized void shutdown() + { + // synchronized, because there is no way to access super.mainLock, which would be + // the preferred way to make this threadsafe + if (!isShutdown()) + unregisterMBean(); + + super.shutdown(); + } + + @Override + public synchronized List shutdownNow() + { + // synchronized, because there is no way to access super.mainLock, which would be + // the preferred way to make this threadsafe + if (!isShutdown()) + unregisterMBean(); + + return super.shutdownNow(); + } + + /** + * Get the number of completed tasks + */ + public long getCompletedTasks() + { + return getCompletedTaskCount(); + } + + /** + * Get the number of tasks waiting to be executed + */ + public long getPendingTasks() + { + return getTaskCount() - getCompletedTaskCount(); + } + + public int getTotalBlockedTasks() + { + return (int) metrics.totalBlocked.count(); + } + + public int getCurrentlyBlockedTasks() + { + return (int) metrics.currentBlocked.count(); + } + + public int getCoreThreads() + { + return getCorePoolSize(); + } + + public void setCoreThreads(int number) + { + setCorePoolSize(number); + } + + public int getMaximumThreads() + { + return getMaximumPoolSize(); + } + + public void setMaximumThreads(int number) + { + setMaximumPoolSize(number); + } +} diff --git a/src/java/org/apache/cassandra/concurrent/JMXEnabledScheduledThreadPoolExecutorMBean.java b/src/java/org/apache/cassandra/concurrent/JMXEnabledScheduledThreadPoolExecutorMBean.java new file mode 100644 index 0000000000..d9c45e3e59 --- /dev/null +++ b/src/java/org/apache/cassandra/concurrent/JMXEnabledScheduledThreadPoolExecutorMBean.java @@ -0,0 +1,26 @@ +/* + * 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.concurrent; + +/** + * @see org.apache.cassandra.metrics.ThreadPoolMetrics + */ +@Deprecated +public interface JMXEnabledScheduledThreadPoolExecutorMBean extends JMXEnabledThreadPoolExecutorMBean +{ +} diff --git a/src/java/org/apache/cassandra/concurrent/JMXEnabledThreadPoolExecutor.java b/src/java/org/apache/cassandra/concurrent/JMXEnabledThreadPoolExecutor.java index 5c96bb6e49..3f60df163b 100644 --- a/src/java/org/apache/cassandra/concurrent/JMXEnabledThreadPoolExecutor.java +++ b/src/java/org/apache/cassandra/concurrent/JMXEnabledThreadPoolExecutor.java @@ -37,7 +37,6 @@ public class JMXEnabledThreadPoolExecutor extends DebuggableThreadPoolExecutor i { private final String mbeanName; private final ThreadPoolMetrics metrics; - public final int maxPoolSize; public JMXEnabledThreadPoolExecutor(String threadPoolName) { @@ -74,7 +73,6 @@ public class JMXEnabledThreadPoolExecutor extends DebuggableThreadPoolExecutor i { super(corePoolSize, maxPoolSize, keepAliveTime, unit, workQueue, threadFactory); super.prestartAllCoreThreads(); - this.maxPoolSize = maxPoolSize; metrics = new ThreadPoolMetrics(this, jmxPath, threadFactory.id); MBeanServer mbs = ManagementFactory.getPlatformMBeanServer(); diff --git a/src/java/org/apache/cassandra/db/HintedHandOffManager.java b/src/java/org/apache/cassandra/db/HintedHandOffManager.java index 8c4477b9cd..2bc26aea41 100644 --- a/src/java/org/apache/cassandra/db/HintedHandOffManager.java +++ b/src/java/org/apache/cassandra/db/HintedHandOffManager.java @@ -34,19 +34,15 @@ import com.google.common.collect.ImmutableSortedSet; import com.google.common.collect.Lists; import com.google.common.util.concurrent.RateLimiter; import com.google.common.util.concurrent.Uninterruptibles; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.concurrent.JMXEnabledThreadPoolExecutor; +import org.apache.cassandra.concurrent.JMXEnabledScheduledThreadPoolExecutor; import org.apache.cassandra.concurrent.NamedThreadFactory; -import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.Schema; -import org.apache.cassandra.db.composites.CellName; -import org.apache.cassandra.db.composites.Composite; -import org.apache.cassandra.db.composites.Composites; import org.apache.cassandra.db.compaction.CompactionManager; +import org.apache.cassandra.db.composites.*; import org.apache.cassandra.db.filter.*; import org.apache.cassandra.db.marshal.Int32Type; import org.apache.cassandra.db.marshal.UUIDType; @@ -63,10 +59,7 @@ import org.apache.cassandra.metrics.HintedHandoffMetrics; import org.apache.cassandra.net.MessageOut; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.service.*; -import org.apache.cassandra.utils.ByteBufferUtil; -import org.apache.cassandra.utils.FBUtilities; -import org.apache.cassandra.utils.JVMStabilityInspector; -import org.apache.cassandra.utils.UUIDGen; +import org.apache.cassandra.utils.*; import org.cliffc.high_scale_lib.NonBlockingHashSet; /** @@ -106,14 +99,13 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean static final int maxHintTTL = Integer.parseInt(System.getProperty("cassandra.maxHintTTL", String.valueOf(Integer.MAX_VALUE))); - private final NonBlockingHashSet queuedDeliveries = new NonBlockingHashSet(); + private final NonBlockingHashSet queuedDeliveries = new NonBlockingHashSet<>(); - private final ThreadPoolExecutor executor = new JMXEnabledThreadPoolExecutor(DatabaseDescriptor.getMaxHintsThread(), - Integer.MAX_VALUE, - TimeUnit.SECONDS, - new LinkedBlockingQueue(), - new NamedThreadFactory("HintedHandoff", Thread.MIN_PRIORITY), - "internal"); + private final JMXEnabledScheduledThreadPoolExecutor executor = + new JMXEnabledScheduledThreadPoolExecutor( + DatabaseDescriptor.getMaxHintsThread(), + new NamedThreadFactory("HintedHandoff", Thread.MIN_PRIORITY), + "internal"); private final ColumnFamilyStore hintStore = Keyspace.open(SystemKeyspace.NAME).getColumnFamilyStore(SystemKeyspace.HINTS); @@ -176,7 +168,7 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean metrics.log(); } }; - ScheduledExecutors.optionalTasks.scheduleWithFixedDelay(runnable, 10, 10, TimeUnit.MINUTES); + executor.scheduleWithFixedDelay(runnable, 10, 10, TimeUnit.MINUTES); } private static void deleteHint(ByteBuffer tokenBytes, CellName columnName, long timestamp) @@ -228,7 +220,7 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean } } }; - ScheduledExecutors.optionalTasks.submit(runnable); + executor.submit(runnable); } //foobar @@ -249,7 +241,7 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean } } }; - ScheduledExecutors.optionalTasks.submit(runnable).get(); + executor.submit(runnable).get(); } @@ -513,7 +505,7 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean IPartitioner p = StorageService.getPartitioner(); RowPosition minPos = p.getMinimumToken().minKeyBound(); - Range range = new Range(minPos, minPos); + Range range = new Range<>(minPos, minPos); IDiskAtomFilter filter = new NamesQueryFilter(ImmutableSortedSet.of()); List rows = hintStore.getRangeSlice(range, null, filter, Integer.MAX_VALUE, System.currentTimeMillis()); for (Row row : rows) @@ -577,7 +569,7 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean Token.TokenFactory tokenFactory = StorageService.getPartitioner().getTokenFactory(); // Extract the keys as strings to be reported. - LinkedList result = new LinkedList(); + LinkedList result = new LinkedList<>(); for (Row row : getHintsSlice(1)) { if (row.cf != null) //ignore removed rows @@ -596,7 +588,7 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean // From keys "" to ""... IPartitioner partitioner = StorageService.getPartitioner(); RowPosition minPos = partitioner.getMinimumToken().minKeyBound(); - Range range = new Range(minPos, minPos); + Range range = new Range<>(minPos, minPos); try { diff --git a/src/java/org/apache/cassandra/metrics/ThreadPoolMetrics.java b/src/java/org/apache/cassandra/metrics/ThreadPoolMetrics.java index 8600e0c184..a5e6dafbb3 100644 --- a/src/java/org/apache/cassandra/metrics/ThreadPoolMetrics.java +++ b/src/java/org/apache/cassandra/metrics/ThreadPoolMetrics.java @@ -19,8 +19,6 @@ package org.apache.cassandra.metrics; import java.util.concurrent.ThreadPoolExecutor; -import org.apache.cassandra.concurrent.JMXEnabledThreadPoolExecutor; - import com.yammer.metrics.Metrics; import com.yammer.metrics.core.*; @@ -54,7 +52,7 @@ public class ThreadPoolMetrics * @param path Type of thread pool * @param poolName Name of thread pool to identify metrics */ - public ThreadPoolMetrics(final JMXEnabledThreadPoolExecutor executor, String path, String poolName) + public ThreadPoolMetrics(final ThreadPoolExecutor executor, String path, String poolName) { this.factory = new ThreadPoolMetricNameFactory("ThreadPools", path, poolName); @@ -85,7 +83,7 @@ public class ThreadPoolMetrics { public Integer value() { - return executor.maxPoolSize; + return executor.getMaximumPoolSize(); } }); }