diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java index 07219abcc5563..f6dd28ede3966 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java @@ -19,6 +19,7 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutorService; +import java.util.concurrent.atomic.LongAdder; import org.apache.camel.AsyncCallback; import org.apache.camel.CamelContext; @@ -30,9 +31,11 @@ import org.apache.camel.Predicate; import org.apache.camel.Processor; import org.apache.camel.Route; +import org.apache.camel.ShutdownRunningTask; import org.apache.camel.Traceable; import org.apache.camel.spi.IdAware; import org.apache.camel.spi.RouteIdAware; +import org.apache.camel.spi.ShutdownAware; import org.apache.camel.spi.StepIdAware; import org.apache.camel.spi.SynchronizationRouteAware; import org.apache.camel.support.ExchangeHelper; @@ -47,7 +50,8 @@ /** * Processor implementing onCompletion. */ -public class OnCompletionProcessor extends BaseProcessorSupport implements Traceable, IdAware, RouteIdAware, StepIdAware { +public class OnCompletionProcessor extends BaseProcessorSupport + implements Traceable, ShutdownAware, IdAware, RouteIdAware, StepIdAware { private static final Logger LOG = LoggerFactory.getLogger(OnCompletionProcessor.class); @@ -64,6 +68,7 @@ public class OnCompletionProcessor extends BaseProcessorSupport implements Trace private final boolean useOriginalBody; private final boolean afterConsumer; private final boolean routeScoped; + private final LongAdder taskCount = new LongAdder(); public OnCompletionProcessor(CamelContext camelContext, Processor processor, ExecutorService executorService, boolean shutdownExecutorService, @@ -107,7 +112,11 @@ protected void doStop() throws Exception { protected void doShutdown() throws Exception { ServiceHelper.stopAndShutdownService(processor); if (shutdownExecutorService) { - getCamelContext().getExecutorServiceManager().shutdownNow(executorService); + List dropped = getCamelContext().getExecutorServiceManager().shutdownNow(executorService); + if (dropped != null && !dropped.isEmpty()) { + // the tasks still queued in the thread pool will never run, so they are no longer pending + taskCount.add(-dropped.size()); + } } } @@ -115,6 +124,22 @@ public CamelContext getCamelContext() { return camelContext; } + @Override + public boolean deferShutdown(ShutdownRunningTask shutdownRunningTask) { + // not in use + return true; + } + + @Override + public int getPendingExchangesSize() { + return taskCount.intValue(); + } + + @Override + public void prepareShutdown(boolean suspendOnly, boolean forced) { + // noop + } + @Override public String getId() { return id; @@ -162,6 +187,30 @@ public boolean process(Exchange exchange, AsyncCallback callback) { return true; } + /** + * Submits the onCompletion task to the thread pool (parallel processing). The task is counted as pending from when + * it is submitted until it is done, so a graceful shutdown waits for it. + */ + @SuppressWarnings("deprecation") + private void submitTask(Runnable task) { + taskCount.increment(); + Runnable counted = () -> { + try { + task.run(); + } finally { + taskCount.decrement(); + } + }; + try { + // Deprecated since 4.19.0 + executorService.submit(prepareMDCParallelTask(camelContext, counted)); + } catch (RuntimeException e) { + // the task will not run + taskCount.decrement(); + throw e; + } + } + protected boolean isCreateCopy() { // we need to create a correlated copy if we run in parallel mode or is in after consumer mode (as the UoW would be done on the original exchange otherwise) return executorService != null || afterConsumer; @@ -301,7 +350,6 @@ public void onAfterRoute(Route route, Exchange exchange) { }; } - @SuppressWarnings("deprecation") @Override public void onComplete(final Exchange exchange) { if (shouldSkip(exchange, onFailureOnly)) { @@ -316,9 +364,7 @@ public void onComplete(final Exchange exchange) { LOG.debug("Processing onComplete: {}", copy); doProcess(processor, copy); }; - // Deprecated since 4.19.0 - task = prepareMDCParallelTask(camelContext, task); - executorService.submit(task); + submitTask(task); } else { // run without thread-pool LOG.debug("Processing onComplete: {}", copy); @@ -326,7 +372,6 @@ public void onComplete(final Exchange exchange) { } } - @SuppressWarnings("deprecation") @Override public void onFailure(final Exchange exchange) { if (shouldSkip(exchange, onCompleteOnly)) { @@ -349,9 +394,7 @@ public void onFailure(final Exchange exchange) { // restore exception after processing copy.setException(original); }; - // Deprecated since 4.19.0 - task = prepareMDCParallelTask(camelContext, task); - executorService.submit(task); + submitTask(task); } else { // run without thread-pool LOG.debug("Processing onFailure: {}", copy); @@ -444,7 +487,6 @@ public void onBeforeRoute(Route route, Exchange exchange) { // NO-OP } - @SuppressWarnings("deprecation") @Override public void onAfterRoute(Route route, Exchange exchange) { LOG.debug("onAfterRoute from Route {}", route.getRouteId()); @@ -479,9 +521,7 @@ public void onAfterRoute(Route route, Exchange exchange) { LOG.debug("Processing onAfterRoute: {}", copy); doProcess(processor, copy); }; - // Deprecated since 4.19.0 - task = prepareMDCParallelTask(camelContext, task); - executorService.submit(task); + submitTask(task); } else { // run without thread-pool LOG.debug("Processing onAfterRoute: {}", copy); diff --git a/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingForcedShutdownTest.java b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingForcedShutdownTest.java new file mode 100644 index 0000000000000..3b6926d8093ee --- /dev/null +++ b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingForcedShutdownTest.java @@ -0,0 +1,87 @@ +/* + * 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.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.CamelContext; +import org.apache.camel.ContextTestSupport; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.builder.ThreadPoolProfileBuilder; +import org.junit.jupiter.api.Test; + +import static org.awaitility.Awaitility.await; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * When a graceful shutdown times out, the thread pool of a parallel onCompletion is shut down, which drops the + * onCompletion tasks still queued in it. They must no longer be counted as pending exchanges. + */ +class OnCompletionParallelProcessingForcedShutdownTest extends ContextTestSupport { + + private final CountDownLatch release = new CountDownLatch(1); + private final AtomicInteger started = new AtomicInteger(); + + @Test + void testForcedShutdownDropsQueuedOnCompletion() throws Exception { + // the first onCompletion runs (and waits), the second is queued in the pool with one thread + template.sendBody("direct:start", "A"); + template.sendBody("direct:start", "B"); + await().atMost(10, TimeUnit.SECONDS).until(() -> started.get() == 1); + OnCompletionProcessor onCompletion = context.getProcessor("oc", OnCompletionProcessor.class); + assertEquals(2, onCompletion.getPendingExchangesSize()); + + // the graceful shutdown times out, and the pool is shut down: the running task is interrupted, and the queued + // task is dropped + context.getShutdownStrategy().setTimeout(1); + context.stop(); + assertTrue(context.getShutdownStrategy().hasTimeoutOccurred()); + + await().atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(0, onCompletion.getPendingExchangesSize(), + "No onCompletion task should be pending after the thread pool is shut down")); + assertEquals(1, started.get(), "The queued onCompletion should not have run"); + } + + @Override + protected CamelContext createCamelContext() throws Exception { + CamelContext context = super.createCamelContext(); + context.getExecutorServiceManager() + .registerThreadPoolProfile(new ThreadPoolProfileBuilder("oneThread").poolSize(1).maxPoolSize(1).build()); + return context; + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("direct:start") + .onCompletion().id("oc").parallelProcessing().executorService("oneThread") + .process(e -> { + started.incrementAndGet(); + release.await(20, TimeUnit.SECONDS); + }) + .end() + .to("mock:result"); + } + }; + } +} diff --git a/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingShutdownTest.java b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingShutdownTest.java new file mode 100644 index 0000000000000..3f69bd85ac21e --- /dev/null +++ b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionParallelProcessingShutdownTest.java @@ -0,0 +1,90 @@ +/* + * 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.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.ContextTestSupport; +import org.apache.camel.builder.RouteBuilder; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import static org.awaitility.Awaitility.await; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * The onCompletion tasks running in the thread pool of a parallel onCompletion are pending exchanges, so a graceful + * shutdown waits for them. + */ +class OnCompletionParallelProcessingShutdownTest extends ContextTestSupport { + + private final CountDownLatch onCompletionStarted = new CountDownLatch(1); + private final CountDownLatch camelStopping = new CountDownLatch(1); + private final AtomicInteger onCompletionDone = new AtomicInteger(); + private final ExecutorService stopper = Executors.newSingleThreadExecutor(); + + @AfterEach + void shutdownStopper() { + stopper.shutdownNow(); + } + + @Test + void testGracefulShutdownWaitsForOnCompletion() throws Exception { + template.sendBody("direct:start", "Hello World"); + assertTrue(onCompletionStarted.await(10, TimeUnit.SECONDS)); + OnCompletionProcessor onCompletion = context.getProcessor("oc", OnCompletionProcessor.class); + assertEquals(1, onCompletion.getPendingExchangesSize(), "The running onCompletion task should be pending"); + + // the exchange is done, and Camel is stopped while its onCompletion is still running + Future stop = stopper.submit(() -> { + context.stop(); + return null; + }); + await().atMost(10, TimeUnit.SECONDS).until(() -> context.isStopping() || context.isStopped()); + camelStopping.countDown(); + stop.get(30, TimeUnit.SECONDS); + + assertEquals(1, onCompletionDone.get(), "The onCompletion should be done"); + assertEquals(0, onCompletion.getPendingExchangesSize(), "No onCompletion task should be pending"); + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("direct:start") + .onCompletion().id("oc").parallelProcessing() + .process(e -> { + onCompletionStarted.countDown(); + // the onCompletion takes until Camel is stopping + if (camelStopping.await(10, TimeUnit.SECONDS)) { + onCompletionDone.incrementAndGet(); + } + }) + .end() + .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 f70b258eedc10..9f98bbeba24e5 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 @@ -782,6 +782,17 @@ seen, and a repository backed by a remote store keeps its connection while the r forget the ids, clear the repository with `IdempotentRepository.clear()` (or the `clear` JMX operation of the Idempotent Consumer). +=== camel-core - a graceful shutdown waits for the parallel onCompletion tasks + +With `onCompletion().parallelProcessing()`, stopping or suspending a route (and stopping `CamelContext`) now waits +for the onCompletion tasks that are running or queued in its thread pool, up to the shutdown timeout, just as it +waits for the inflight exchanges. Previously the shutdown did not wait for them, and when the thread pool was shut +down, the queued tasks were dropped and the running ones were interrupted. + +An onCompletion with `parallelProcessing` that synchronously stops its own route now waits for itself until the +shutdown timeout occurs, and the route is then stopped forcibly. Stop the route asynchronously instead, for example +from a separate thread or with the Control Bus `async=true` option. + === Component deprecation ==== camel-minio