This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 3468c7edf253 CAMEL-25073: camel-management - Give each thread pool of
the same source its own MBean (item 4) (#27465)
3468c7edf253 is described below
commit 3468c7edf2531e0e9cad35495fe77d040116703e
Author: allthingssecurity <[email protected]>
AuthorDate: Wed Oct 7 13:40:06 2026 +0530
CAMEL-25073: camel-management - Give each thread pool of the same source
its own MBean (item 4) (#27465)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../org/apache/camel/spi/LifecycleStrategy.java | 8 +-
.../impl/engine/BaseExecutorServiceManager.java | 13 ++-
.../ManagedAggregateThreadPoolsTest.java | 96 ++++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 10 +++
4 files changed, 121 insertions(+), 6 deletions(-)
diff --git
a/core/camel-api/src/main/java/org/apache/camel/spi/LifecycleStrategy.java
b/core/camel-api/src/main/java/org/apache/camel/spi/LifecycleStrategy.java
index 32553c954f1c..f72615072e48 100644
--- a/core/camel-api/src/main/java/org/apache/camel/spi/LifecycleStrategy.java
+++ b/core/camel-api/src/main/java/org/apache/camel/spi/LifecycleStrategy.java
@@ -209,7 +209,9 @@ public interface LifecycleStrategy {
* @param camelContext the camel context
* @param threadPool the thread pool
* @param id id of the thread pool
- * @param sourceId id of the source creating the thread pool
(can be null in special cases)
+ * @param sourceId id of the source creating the thread pool,
or the name of the thread pool when the id
+ * of the thread pool is derived from the
source instance, as such a source can create
+ * several thread pools (can be null in special
cases)
* @param routeId id of the route for the source (is null if
no source)
* @param threadPoolProfileId id of the thread pool profile, if used for
creating this thread pool (can be null)
*/
@@ -232,7 +234,9 @@ public interface LifecycleStrategy {
* @param camelContext the camel context
* @param executorService the executor service
* @param id id of the thread pool
- * @param sourceId id of the source creating the thread pool
(can be null in special cases)
+ * @param sourceId id of the source creating the thread pool,
or the name of the thread pool when the id
+ * of the thread pool is derived from the
source instance, as such a source can create
+ * several thread pools (can be null in special
cases)
* @param routeId id of the route for the source (is null if
no source)
* @param threadPoolProfileId id of the thread pool profile, if used for
creating this thread pool (can be null)
*/
diff --git
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/BaseExecutorServiceManager.java
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/BaseExecutorServiceManager.java
index 249777abe40c..78a0a0c28198 100644
---
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/BaseExecutorServiceManager.java
+++
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/BaseExecutorServiceManager.java
@@ -202,7 +202,7 @@ public class BaseExecutorServiceManager extends
ServiceSupport implements Execut
ThreadFactory threadFactory = createThreadFactory(source,
sanitizedName, true);
ExecutorService executorService =
threadPoolFactory.newThreadPool(profile, threadFactory);
- onThreadPoolCreated(executorService, source, profile.getId());
+ onThreadPoolCreated(executorService, source, sanitizedName,
profile.getId());
if (LOG.isDebugEnabled()) {
LOG.debug("Created new ThreadPool for source: {} with name: {}. ->
{}", source, sanitizedName, executorService);
}
@@ -227,7 +227,7 @@ public class BaseExecutorServiceManager extends
ServiceSupport implements Execut
public ExecutorService newCachedThreadPool(Object source, String name) {
String sanitizedName = URISupport.sanitizeUri(name);
ExecutorService answer =
threadPoolFactory.newCachedThreadPool(createThreadFactory(source,
sanitizedName, true));
- onThreadPoolCreated(answer, source, null);
+ onThreadPoolCreated(answer, source, sanitizedName, null);
if (LOG.isDebugEnabled()) {
LOG.debug("Created new CachedThreadPool for source: {} with name:
{}. -> {}", source, sanitizedName, answer);
@@ -261,7 +261,7 @@ public class BaseExecutorServiceManager extends
ServiceSupport implements Execut
profile.addDefaults(getDefaultThreadPoolProfile());
ScheduledExecutorService answer
= threadPoolFactory.newScheduledThreadPool(profile,
createThreadFactory(source, sanitizedName, true));
- onThreadPoolCreated(answer, source, null);
+ onThreadPoolCreated(answer, source, sanitizedName, null);
if (LOG.isDebugEnabled()) {
LOG.debug("Created new ScheduledThreadPool for source: {} with
name: {} -> {}", source, sanitizedName, answer);
@@ -550,9 +550,11 @@ public class BaseExecutorServiceManager extends
ServiceSupport implements Execut
*
* @param executorService the thread pool
* @param source the source to use the thread pool
+ * @param name the name of the thread pool
* @param threadPoolProfileId profile id, if the thread pool was created
from a thread pool profile
*/
- private void onThreadPoolCreated(ExecutorService executorService, Object
source, String threadPoolProfileId) {
+ private void onThreadPoolCreated(
+ ExecutorService executorService, Object source, String name,
String threadPoolProfileId) {
// add to internal list of thread pools
executorServices.add(executorService);
@@ -574,6 +576,9 @@ public class BaseExecutorServiceManager extends
ServiceSupport implements Execut
} else {
// fallback and use the simple class name with hashcode for
the id so its unique for this given source
id = source.getClass().getSimpleName() + "(" +
ObjectHelper.getIdentityHashCode(source) + ")";
+ // and the name of the thread pool as source id, as a source
can create several thread pools
+ // (such as the timeout checker and the optimistic locking
executor of the aggregator)
+ sourceId = name;
}
} else {
// no source, so fallback and use the simple class name from
thread pool and its hashcode identity so its unique
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedAggregateThreadPoolsTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedAggregateThreadPoolsTest.java
new file mode 100644
index 000000000000..9f9e0ca785f3
--- /dev/null
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedAggregateThreadPoolsTest.java
@@ -0,0 +1,96 @@
+/*
+ * 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.camel.management;
+
+import java.util.List;
+import java.util.Set;
+
+import javax.management.MBeanServer;
+import javax.management.ObjectName;
+
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.processor.aggregate.AggregateProcessor;
+import org.apache.camel.processor.aggregate.UseLatestAggregationStrategy;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.condition.DisabledOnOs;
+import org.junit.jupiter.api.condition.OS;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * The aggregator creates several thread pools from the same source (the
AggregateProcessor), which must each have their
+ * own MBean.
+ */
+@DisabledOnOs(OS.AIX)
+class ManagedAggregateThreadPoolsTest extends ManagementTestSupport {
+
+ @Test
+ void testThreadPoolsOfTheSameSource() throws Exception {
+ MBeanServer mbeanServer = getMBeanServer();
+
+ List<ObjectName> pools = aggregateProcessorPools(mbeanServer);
+ assertEquals(2, pools.size(), "The timeout checker and the optimistic
locking pools should be registered: " + pools);
+
+ List<String> sourceIds = pools.stream().map(on ->
getSourceId(mbeanServer, on)).sorted().toList();
+
assertEquals(List.of(AggregateProcessor.AGGREGATE_OPTIMISTIC_LOCKING_EXECUTOR,
+ AggregateProcessor.AGGREGATE_TIMEOUT_CHECKER), sourceIds);
+ for (ObjectName on : pools) {
+ String name = ObjectName.unquote(on.getKeyProperty("name"));
+ String id = (String) mbeanServer.getAttribute(on, "Id");
+ assertTrue(id.startsWith("AggregateProcessor("), id);
+ assertEquals(id + "(" + getSourceId(mbeanServer, on) + ")", name);
+ }
+
+ // and after a restart of the route
+ context.getRouteController().stopRoute("foo");
+ context.getRouteController().startRoute("foo");
+ assertEquals(sourceIds,
aggregateProcessorPools(mbeanServer).stream().map(on ->
getSourceId(mbeanServer, on))
+ .sorted().toList());
+
+ // both are unregistered when the route is removed
+ context.getRouteController().stopRoute("foo");
+ context.removeRoute("foo");
+ assertEquals(List.of(), aggregateProcessorPools(mbeanServer));
+ }
+
+ private static List<ObjectName> aggregateProcessorPools(MBeanServer
mbeanServer) throws Exception {
+ Set<ObjectName> pools = mbeanServer.queryNames(new
ObjectName("*:type=threadpools,*"), null);
+ return pools.stream().filter(on ->
on.getKeyProperty("name").startsWith("\"AggregateProcessor(")).toList();
+ }
+
+ private static String getSourceId(MBeanServer mbeanServer, ObjectName on) {
+ try {
+ return (String) mbeanServer.getAttribute(on, "SourceId");
+ } catch (Exception e) {
+ throw new AssertionError(e);
+ }
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:foo").routeId("foo")
+ .aggregate(constant(true), new
UseLatestAggregationStrategy()).completionTimeout(1000)
+ .optimisticLocking()
+ .to("mock:result");
+ }
+ };
+ }
+}
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index bfb4845a5b4d..9c7d82c1a1f7 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -700,6 +700,16 @@ then redelivered.
=== camel-management
+The thread pool MBeans of a source that creates several thread pools (such as
the timeout checker and the optimistic
+locking pools of the Aggregate EIP) had the same name, so only the first one
was registered. For a source that is not
+an EIP, a name or a static service, the `SourceId` attribute is now the name
of the thread pool, which is also added to
+the MBean name: `AggregateProcessor(0x1b2c3d4e)(AggregateTimeoutChecker)`
instead of `AggregateProcessor(0x1b2c3d4e)`.
+This also applies to the thread pools of component consumers, producers and
their helper services, such as the
+consumer threads of `seda` (`SedaConsumer(0x1b2c3d4e)(seda://foo)`), the reply
manager pools of request/reply over
+`jms` and `sjms`
(`JmsProducer(0x1b2c3d4e)(JmsReplyManagerTimeoutChecker[bar])`), the timeout
pools of `netty`
+request/reply and the read lock release task of `file`.
+The names of the thread pools of EIPs, of pools created by name and of static
services do not change.
+
The JMX `browse` operation of the `DefaultInflightRepository` MBean and the
`listAwaitThreads` data of the
`DefaultAsyncProcessorAwaitManager` MBean gained a `nodeSource` column, saying
where the node is in the source (such
as `orders.camel.yaml:18`), next to the existing `nodeId`. It is `null` when
message history or source location is