This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch backport/25036-to-camel-4.18.x in repository https://gitbox.apache.org/repos/asf/camel.git
commit af7c70711420fb7009ae76e3d041f6a5968d150c Author: Mark Snijder <[email protected]> AuthorDate: Thu Jul 23 17:11:33 2026 +0200 CAMEL-24248: Stop services released after pool eviction in ServicePool When pollEnrich processes concurrent exchanges using different dynamic endpoint URIs and its consumer cache is smaller than the concurrency level, a polling consumer's pool can be evicted while the consumer is still in use. The released consumer was never stopped, causing resource leaks and steadily increasing memory usage. Add cleanup in ServicePool.release() to stop and remove a service when its endpoint pool has already been evicted. Extract common stop-and-remove logic into a shared method. Closes #25036 Co-authored-by: GPT-5.6 Sol <[email protected]> --- .../PollEnrichConcurrentConsumersCacheTest.java | 239 +++++++++++++++++++++ .../apache/camel/support/cache/ServicePool.java | 19 +- 2 files changed, 252 insertions(+), 6 deletions(-) diff --git a/core/camel-core/src/test/java/org/apache/camel/processor/PollEnrichConcurrentConsumersCacheTest.java b/core/camel-core/src/test/java/org/apache/camel/processor/PollEnrichConcurrentConsumersCacheTest.java new file mode 100644 index 000000000000..ff43aeef736c --- /dev/null +++ b/core/camel-core/src/test/java/org/apache/camel/processor/PollEnrichConcurrentConsumersCacheTest.java @@ -0,0 +1,239 @@ +/* + * 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.processor; + +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.Consumer; +import org.apache.camel.ContextTestSupport; +import org.apache.camel.Endpoint; +import org.apache.camel.Exchange; +import org.apache.camel.PollingConsumer; +import org.apache.camel.Processor; +import org.apache.camel.Producer; +import org.apache.camel.RuntimeCamelException; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.support.DefaultComponent; +import org.apache.camel.support.DefaultEndpoint; +import org.apache.camel.support.PollingConsumerSupport; +import org.apache.camel.support.service.ServiceHelper; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +class PollEnrichConcurrentConsumersCacheTest extends ContextTestSupport { + private static final long TIMEOUT_SECONDS = 10; + + private final Map<String, PollState> states = new ConcurrentHashMap<>(); + private final List<PollingConsumer> consumers = new CopyOnWriteArrayList<>(); + + @AfterEach + void cleanUpConsumers() { + states.values().forEach(PollState::release); + ServiceHelper.stopAndShutdownServices(consumers); + } + + @Test + void concurrentDynamicPollsStopEvictedConsumer() throws Exception { + PollState a = registerState("a", true); + PollState b = registerState("b", false); + context.addComponent("cachetest", pollingComponent()); + + Future<String> resultA + = template.asyncRequestBodyAndHeader("seda:start", "trigger-a", "pollUri", "cachetest:a", String.class); + assertThat(a.getEnteredReceive().await(TIMEOUT_SECONDS, TimeUnit.SECONDS)) + .as("consumer A entered receive") + .isTrue(); + + Future<String> resultB + = template.asyncRequestBodyAndHeader("seda:start", "trigger-b", "pollUri", "cachetest:b", String.class); + assertThat(b.getEnteredReceive().await(TIMEOUT_SECONDS, TimeUnit.SECONDS)) + .as("consumer B entered receive") + .isTrue(); + + assertThat(resultB.get(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isEqualTo("b"); + a.release(); + assertThat(resultA.get(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isEqualTo("a"); + + assertThat(a.getCreated()).hasValue(1); + assertThat(a.getReceived()).hasValue(1); + assertThat(b.getCreated()).hasValue(1); + assertThat(b.getReceived()).hasValue(1); + assertThat(b.getStopped()).as("cached consumer B remains started").hasValue(0); + assertThat(a.getStopped()).as("evicted consumer A is stopped when it is released").hasValue(1); + } + + @Test + void sequentialDynamicPollsStopIdleConsumer() throws Exception { + PollState a = registerState("a", false); + PollState b = registerState("b", false); + context.addComponent("cachetest", pollingComponent()); + + String resultA = template.requestBodyAndHeader("seda:start", "trigger-a", "pollUri", "cachetest:a", String.class); + + assertThat(resultA).isEqualTo("a"); + assertThat(a.getCreated()).hasValue(1); + assertThat(a.getReceived()).hasValue(1); + assertThat(a.getStopped()).as("consumer A is idle in its pool").hasValue(0); + + String resultB = template.requestBodyAndHeader("seda:start", "trigger-b", "pollUri", "cachetest:b", String.class); + + assertThat(resultB).isEqualTo("b"); + assertThat(b.getCreated()).hasValue(1); + assertThat(b.getReceived()).hasValue(1); + assertThat(a.getStopped()).as("idle consumer A is stopped during eviction").hasValue(1); + assertThat(b.getStopped()).as("cached consumer B remains started").hasValue(0); + } + + private PollState registerState(String name, boolean blockReceive) { + PollState state = new PollState(blockReceive); + states.put(name, state); + return state; + } + + private DefaultComponent pollingComponent() { + return new DefaultComponent() { + @Override + protected Endpoint createEndpoint(String uri, String remaining, Map<String, Object> parameters) { + PollState state = states.get(remaining); + if (state == null) { + throw new IllegalArgumentException("No polling state registered for " + remaining); + } + + return new DefaultEndpoint(uri, this) { + @Override + public boolean isSingletonProducer() { + return false; + } + + @Override + public Producer createProducer() { + return null; + } + + @Override + public Consumer createConsumer(Processor processor) { + return null; + } + + @Override + public PollingConsumer createPollingConsumer() { + state.getCreated().incrementAndGet(); + PollingConsumerSupport consumer = new PollingConsumerSupport(this) { + @Override + public Exchange receive() { + return poll(); + } + + @Override + public Exchange receive(long timeout) { + return poll(); + } + + @Override + public Exchange receiveNoWait() { + return poll(); + } + + private Exchange poll() { + state.getEnteredReceive().countDown(); + try { + if (!state.getAllowReceiveToFinish().await(TIMEOUT_SECONDS, TimeUnit.SECONDS)) { + throw new IllegalStateException("Timed out waiting to finish polling " + remaining); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeCamelException(e); + } + + state.getReceived().incrementAndGet(); + Exchange exchange = getEndpoint().createExchange(); + exchange.getMessage().setBody(remaining); + return exchange; + } + + @Override + protected void doStop() { + state.getStopped().incrementAndGet(); + } + }; + consumers.add(consumer); + return consumer; + } + }; + } + }; + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("seda:start?concurrentConsumers=4") + .pollEnrich() + .header("pollUri") + .timeout(5000) + .cacheSize(1); + } + }; + } + + private static final class PollState { + private final AtomicInteger created = new AtomicInteger(); + private final AtomicInteger received = new AtomicInteger(); + private final AtomicInteger stopped = new AtomicInteger(); + private final CountDownLatch enteredReceive = new CountDownLatch(1); + private final CountDownLatch allowReceiveToFinish; + + private PollState(boolean blockReceive) { + this.allowReceiveToFinish = new CountDownLatch(blockReceive ? 1 : 0); + } + + private void release() { + allowReceiveToFinish.countDown(); + } + + public AtomicInteger getCreated() { + return created; + } + + public AtomicInteger getReceived() { + return received; + } + + public AtomicInteger getStopped() { + return stopped; + } + + public CountDownLatch getEnteredReceive() { + return enteredReceive; + } + + public CountDownLatch getAllowReceiveToFinish() { + return allowReceiveToFinish; + } + } +} diff --git a/core/camel-support/src/main/java/org/apache/camel/support/cache/ServicePool.java b/core/camel-support/src/main/java/org/apache/camel/support/cache/ServicePool.java index 91a072a273e9..39853362e907 100644 --- a/core/camel-support/src/main/java/org/apache/camel/support/cache/ServicePool.java +++ b/core/camel-support/src/main/java/org/apache/camel/support/cache/ServicePool.java @@ -94,12 +94,7 @@ abstract class ServicePool<S extends Service> extends ServiceSupport implements } } else { // service no longer in a pool (such as being released twice, or can happen during shutdown of Camel etc) - ServicePool.stop(s); - try { - e.getCamelContext().removeService(s); - } catch (Exception ex) { - LOG.debug("Error removing service: {}. This exception is ignored.", s, ex); - } + stopAndRemove(s); } } @@ -130,6 +125,18 @@ abstract class ServicePool<S extends Service> extends ServiceSupport implements Pool<S> p = pool.get(endpoint); if (p != null) { p.release(s); + } else { + // the pool was evicted while this service was in use + stopAndRemove(s); + } + } + + private void stopAndRemove(S s) { + ServicePool.stop(s); + try { + getEndpoint.apply(s).getCamelContext().removeService(s); + } catch (Exception e) { + LOG.debug("Error removing service: {}. This exception is ignored.", s, e); } }
