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 4b9ece08c9dd CAMEL-24146: Fix FluentProducerTemplate leak in
InMemorySagaCoordinator
4b9ece08c9dd is described below
commit 4b9ece08c9dd8945648c772bb4c57aa47e5a3af8
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 14:09:00 2026 +0200
CAMEL-24146: Fix FluentProducerTemplate leak in InMemorySagaCoordinator
Co-Authored-By: Claude Opus 4.6 <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
Co-authored-by: Claude Opus 4.6 <[email protected]>
---
.../org/apache/camel/saga/InMemorySagaCoordinator.java | 2 +-
.../java/org/apache/camel/saga/InMemorySagaService.java | 14 ++++++++++++++
2 files changed, 15 insertions(+), 1 deletion(-)
diff --git
a/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaCoordinator.java
b/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaCoordinator.java
index 3d85357076d9..f6996eda5cfc 100644
---
a/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaCoordinator.java
+++
b/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaCoordinator.java
@@ -210,7 +210,7 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
Exchange target = createExchange(exchange, endpoint, step);
return CompletableFuture.supplyAsync(() -> {
- Exchange res =
camelContext.createFluentProducerTemplate().to(endpoint).withExchange(target).send();
+ Exchange res = sagaService.getProducerTemplate().send(endpoint,
target);
Exception ex = res.getException();
if (ex != null) {
throw new RuntimeCamelException(res.getException());
diff --git
a/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaService.java
b/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaService.java
index 1f90d61a21e7..11b55c387a82 100644
---
a/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaService.java
+++
b/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaService.java
@@ -23,6 +23,7 @@ import java.util.concurrent.ScheduledExecutorService;
import org.apache.camel.CamelContext;
import org.apache.camel.Exchange;
+import org.apache.camel.ProducerTemplate;
import org.apache.camel.support.service.ServiceSupport;
import org.apache.camel.util.ObjectHelper;
@@ -41,6 +42,8 @@ public class InMemorySagaService extends ServiceSupport
implements CamelSagaServ
private ScheduledExecutorService executorService;
+ private ProducerTemplate producerTemplate;
+
private int maxRetryAttempts = DEFAULT_MAX_RETRY_ATTEMPTS;
private long retryDelayInMilliseconds =
DEFAULT_RETRY_DELAY_IN_MILLISECONDS;
@@ -76,10 +79,17 @@ public class InMemorySagaService extends ServiceSupport
implements CamelSagaServ
this.executorService = camelContext.getExecutorServiceManager()
.newDefaultScheduledThreadPool(this, "saga");
}
+ if (this.producerTemplate == null) {
+ this.producerTemplate = camelContext.createProducerTemplate();
+ }
}
@Override
protected void doStop() throws Exception {
+ if (this.producerTemplate != null) {
+ this.producerTemplate.stop();
+ this.producerTemplate = null;
+ }
if (this.executorService != null) {
camelContext.getExecutorServiceManager().shutdownGraceful(this.executorService);
this.executorService = null;
@@ -90,6 +100,10 @@ public class InMemorySagaService extends ServiceSupport
implements CamelSagaServ
return executorService;
}
+ public ProducerTemplate getProducerTemplate() {
+ return producerTemplate;
+ }
+
@Override
public void setCamelContext(CamelContext camelContext) {
this.camelContext = camelContext;