From ddfdf5d69f4a58193cd8cde3bf76aed03d5c08eb Mon Sep 17 00:00:00 2001 From: Dmitry Konstantinov Date: Sat, 20 Jun 2026 16:49:16 +0100 Subject: [PATCH] Reduce number of scheduledTasks on metric id release in ThreadLocalMetrics Use a single one-time scheduled task with two tick-tock buffers to recycle metric IDs after a sufficiently long delay. Use phantom references for ThreadLocalMetrics cleanup ony if it is needed to reduce the references processing overhead. patch by Dmitry Konstantinov; reviewed by Benedict Elliott Smith for CASSANDRA-21475 --- CHANGES.txt | 1 + .../cassandra/metrics/ThreadLocalMetrics.java | 115 ++++++++++++++---- .../metrics/FreeMetricIdSetTrackerTest.java | 103 ++++++++++++++++ .../metrics/ThreadLocalCounterTest.java | 2 +- 4 files changed, 196 insertions(+), 25 deletions(-) create mode 100644 test/unit/org/apache/cassandra/metrics/FreeMetricIdSetTrackerTest.java diff --git a/CHANGES.txt b/CHANGES.txt index eb29172879..b257a758e0 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 6.0-alpha2 + * Reduce number of scheduledTasks on metric id release in ThreadLocalMetrics (CASSANDRA-21475) * Cache various Enum.values() used in deserialization to avoid per-read array allocation (CASSANDRA-21528) * Fix Accord transaction error message when altering a table (CASSANDRA-20580) * Depend only on platform-specific Zstd JNI native libraries (CASSANDRA-21483) diff --git a/src/java/org/apache/cassandra/metrics/ThreadLocalMetrics.java b/src/java/org/apache/cassandra/metrics/ThreadLocalMetrics.java index feb96c02bd..d7b429ada6 100644 --- a/src/java/org/apache/cassandra/metrics/ThreadLocalMetrics.java +++ b/src/java/org/apache/cassandra/metrics/ThreadLocalMetrics.java @@ -26,6 +26,8 @@ import java.util.List; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; @@ -41,6 +43,7 @@ import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.concurrent.Shutdownable; import io.netty.util.concurrent.FastThreadLocal; +import io.netty.util.concurrent.FastThreadLocalThread; import static com.google.common.collect.ImmutableList.of; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; @@ -64,10 +67,8 @@ public class ThreadLocalMetrics static final AtomicInteger idGenerator = new AtomicInteger(); - private static final Object freeMetricIdSetGuard = new Object(); - @VisibleForTesting - static final BitSet freeMetricIdSet = new BitSet(); + static final FreeMetricIdSetTracker freeMetricIdSetTracker = new FreeMetricIdSetTracker(); private static final List allThreadLocalMetrics = new CopyOnWriteArrayList<>(); @@ -101,7 +102,13 @@ public class ThreadLocalMetrics { ThreadLocalMetrics result = new ThreadLocalMetrics(); allThreadLocalMetrics.add(result); - destroyWhenUnreachable(Thread.currentThread(), result::release); + + Thread thread = Thread.currentThread(); + // use phantom references ony if needed + // CassandraThread is FastThreadLocalThread too + if (!(thread instanceof FastThreadLocalThread)) + destroyWhenUnreachable(thread, result::release); + return result; } @@ -317,13 +324,7 @@ public class ThreadLocalMetrics static int allocateMetricId() { - int metricId; - synchronized (freeMetricIdSetGuard) - { - metricId = freeMetricIdSet.nextSetBit(0); - if (metricId >= 0) - freeMetricIdSet.clear(metricId); - } + int metricId = freeMetricIdSetTracker.getFreeMetricId(); if (metricId < 0) metricId = idGenerator.getAndIncrement(); @@ -374,26 +375,92 @@ public class ThreadLocalMetrics lock.unlock(); } - // there's no an obvious happens-before relation between currentCounterValues[metricId] = 0 write we just did - // and an initial read of the entry by a thread which updates the reused metric - // as a workaround we introduce a delay in recyling to provide the write visibility in practice - // even if it is not formally guaranteed by the JMM - ScheduledExecutors.scheduledTasks.schedule(() -> { - synchronized (freeMetricIdSetGuard) + freeMetricIdSetTracker.markAsFree(metricId); + } + + @VisibleForTesting + static class FreeMetricIdSetTracker + { + private final BitSet freeMetricIdSet = new BitSet(); + + private final BitSet tickDelayedToFreeMetricIdSet = new BitSet(); + private final BitSet tockDelayedToFreeMetricIdSet = new BitSet(); + + private BitSet delayedToFreeMetricIdSet = tickDelayedToFreeMetricIdSet; + + private ScheduledFuture cleanupTask; + + @VisibleForTesting + synchronized void triggerRecycling() + { + cleanupTask = null; + BitSet toProcess = otherSet(delayedToFreeMetricIdSet); + freeMetricIdSet.or(toProcess); + toProcess.clear(); + if (!delayedToFreeMetricIdSet.isEmpty()) + scheduleCleanupTask(); + delayedToFreeMetricIdSet = toProcess; + } + + private BitSet otherSet(BitSet set) + { + return set == tickDelayedToFreeMetricIdSet ? tockDelayedToFreeMetricIdSet : tickDelayedToFreeMetricIdSet; + } + + public synchronized int getFreeMetricId() + { + int metricId = freeMetricIdSet.nextSetBit(0); + if (metricId >= 0) + freeMetricIdSet.clear(metricId); + return metricId; + } + + public synchronized void markAsFree(int metricId) + { + // there's no an obvious happens-before relation between currentCounterValues[metricId] = 0 write we just did + // and an initial read of the entry by a thread which updates the reused metric + // as a workaround we introduce a delay in recyling to provide the write visibility in practice + // even if it is not formally guaranteed by the JMM + delayedToFreeMetricIdSet.set(metricId); + scheduleCleanupTask(); + } + + // must be called while holding this monitor (from a synchronized method) + @VisibleForTesting + protected void scheduleCleanupTask() + { + try { - freeMetricIdSet.set(metricId); + if (cleanupTask == null) + cleanupTask = ScheduledExecutors.scheduledTasks.schedule(this::triggerRecycling, 5, TimeUnit.SECONDS); } - }, 5, TimeUnit.SECONDS); + catch (RejectedExecutionException e) + { + // ignore theoretically possible rejections during a shutdown + } + } + + public synchronized int getFreeMetricSetCardinality() + { + return freeMetricIdSet.cardinality(); + } + + @Override + public synchronized String toString() + { + return "FreeMetricIdSetTracker{" + + "freeMetricIdSet=" + freeMetricIdSet + + ", tickDelayedToFreeMetricIdSet=" + tickDelayedToFreeMetricIdSet + + ", tockDelayedToFreeMetricIdSet=" + tockDelayedToFreeMetricIdSet + + ", delayedToFreeMetricIdSet=" + (delayedToFreeMetricIdSet == tickDelayedToFreeMetricIdSet ? "tick" : "tock") + + '}'; + } } @VisibleForTesting static int getAllocatedMetricsCount() { - int freeCount; - synchronized (freeMetricIdSetGuard) - { - freeCount = freeMetricIdSet.cardinality(); - } + int freeCount = freeMetricIdSetTracker.getFreeMetricSetCardinality(); return idGenerator.get() - freeCount; } diff --git a/test/unit/org/apache/cassandra/metrics/FreeMetricIdSetTrackerTest.java b/test/unit/org/apache/cassandra/metrics/FreeMetricIdSetTrackerTest.java new file mode 100644 index 0000000000..5b7ad96f95 --- /dev/null +++ b/test/unit/org/apache/cassandra/metrics/FreeMetricIdSetTrackerTest.java @@ -0,0 +1,103 @@ +/* + * 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.metrics; + +import java.util.HashSet; +import java.util.Set; + +import com.google.common.collect.ImmutableSet; + +import org.junit.Test; + +import org.apache.cassandra.metrics.ThreadLocalMetrics.FreeMetricIdSetTracker; + +import static org.junit.Assert.assertEquals; + +public class FreeMetricIdSetTrackerTest +{ + // free and release ids in interleaved portions, and verify that a later portion does not + // become reusable before its own two-cycle delay has elapsed - i.e. only the ids that are + // actually due are released, never the freshly freed ones. + @Test + public void testInterleavedFreeAndReleaseReleasesOnlyExpectedIds() + { + FreeMetricIdSetTracker tracker = getTrackerToTest(); + + // portion A: free, then one release cycle (A is now in the delayed buffer, not yet reusable) + Set portionA = ImmutableSet.of(1, 2, 3); + portionA.forEach(tracker::markAsFree); + tracker.triggerRecycling(); + assertEquals(0, tracker.getFreeMetricSetCardinality()); + + // portion B: free a fresh batch, then another release cycle. + // this cycle is A's second cycle (A becomes reusable) but only B's first (B must stay held). + Set portionB = ImmutableSet.of(10, 11); + portionB.forEach(tracker::markAsFree); + tracker.triggerRecycling(); + + // only portion A is released here - portion B must not leak out yet + assertEquals(portionA.size(), tracker.getFreeMetricSetCardinality()); + assertEquals(portionA, drainFreeIds(tracker)); + + // a further release cycle finally makes portion B reusable, and nothing else + tracker.triggerRecycling(); + assertEquals(portionB.size(), tracker.getFreeMetricSetCardinality()); + assertEquals(portionB, drainFreeIds(tracker)); + + assertEquals(0, tracker.getFreeMetricSetCardinality()); + assertEquals(-1, tracker.getFreeMetricId()); + } + + @Test + public void testTriggerRecyclingIsNoOpWhenNothingFreed() + { + FreeMetricIdSetTracker tracker = getTrackerToTest(); + + // triggering with an empty tracker must not produce phantom free ids + tracker.triggerRecycling(); + tracker.triggerRecycling(); + + assertEquals(0, tracker.getFreeMetricSetCardinality()); + assertEquals(-1, tracker.getFreeMetricId()); + } + + private static FreeMetricIdSetTracker getTrackerToTest() + { + return new FreeMetricIdSetTracker() + { + @Override + protected void scheduleCleanupTask() + { + // disable scheduling to make the test deterministic + } + }; + } + + /** + * Drains every id currently available for reuse out of the tracker. + */ + private static Set drainFreeIds(FreeMetricIdSetTracker tracker) + { + Set ids = new HashSet<>(); + int id; + while ((id = tracker.getFreeMetricId()) >= 0) + ids.add(id); + return ids; + } +} diff --git a/test/unit/org/apache/cassandra/metrics/ThreadLocalCounterTest.java b/test/unit/org/apache/cassandra/metrics/ThreadLocalCounterTest.java index 2597504522..03a0cb0640 100644 --- a/test/unit/org/apache/cassandra/metrics/ThreadLocalCounterTest.java +++ b/test/unit/org/apache/cassandra/metrics/ThreadLocalCounterTest.java @@ -90,7 +90,7 @@ public class ThreadLocalCounterTest LOGGER.info("id generator state: {}, free IDs: {}", ThreadLocalMetrics.idGenerator.get(), - ThreadLocalMetrics.freeMetricIdSet); + ThreadLocalMetrics.freeMetricIdSetTracker); LOGGER.info("iteration completed: {} / {}", iteration + 1, ITERATIONS_COUNT); } }