This is an automated email from the ASF dual-hosted git repository.

gnodet pushed a commit to branch 
investigate-and-fix-camel-24272-use-backgroundtas
in repository https://gitbox.apache.org/repos/asf/camel.git

commit 75bca32e26bb43e9b247ae779385c34d910fda9e
Author: Guillaume Nodet <[email protected]>
AuthorDate: Mon Jul 27 16:26:25 2026 +0200

    CAMEL-24272: Migrate reconnection loops from ForegroundTask to 
BackgroundTask
    
    Switch all reconnection and retry loops from ForegroundTask to
    BackgroundTask so they register with the TaskManagerRegistry and
    become visible in the dev console, CLI and TUI internal-tasks view.
    
    Affected components:
    - camel-kafka: KafkaConsumer create/subscribe retry loops
    - camel-ftp: FTP and SFTP reconnection loops
    - camel-mllp: TCP server bind retry loop
    - camel-infinispan: Schema registration retry loop
    - camel-google-pubsub: Error recovery delay
    - camel-zookeeper: Reconnection backoff delay
    - camel-hazelcast: SEDA error recovery delay
    - camel-salesforce: Handshake and subscribe retry delays
    - camel-mongodb-gridfs: GridFS polling loop
    
    Each migration:
    - Replaces Tasks.foregroundTask() with Tasks.backgroundTask()
    - Switches from IterationBudget to TimeBudget via iterationTimeBudget()
    - Adds a ScheduledExecutorService with proper lifecycle management
    - Assigns a descriptive task name for visibility in monitoring tools
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
 .../camel/component/file/remote/FtpOperations.java | 55 ++++++++++-------
 .../component/file/remote/SftpOperations.java      | 36 +++++++----
 .../google/pubsub/GooglePubsubConsumer.java        | 14 ++++-
 .../hazelcast/seda/HazelcastSedaConsumer.java      | 14 ++++-
 .../infinispan/remote/InfinispanRemoteManager.java | 70 ++++++++++++----------
 .../camel/component/kafka/KafkaFetchRecords.java   | 28 ++++++---
 .../mllp/internal/TcpServerBindThread.java         | 35 +++++++----
 .../component/mongodb/gridfs/GridFsConsumer.java   | 18 ++++--
 .../internal/streaming/SubscriptionHelper.java     | 22 +++++--
 .../component/zookeeper/ZooKeeperConsumer.java     | 14 ++++-
 10 files changed, 204 insertions(+), 102 deletions(-)

diff --git 
a/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/FtpOperations.java
 
b/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/FtpOperations.java
index b25f13eed9c9..1dc597eb7944 100644
--- 
a/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/FtpOperations.java
+++ 
b/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/FtpOperations.java
@@ -25,6 +25,7 @@ import java.io.InputStream;
 import java.io.OutputStream;
 import java.time.Duration;
 import java.util.Iterator;
+import java.util.concurrent.ScheduledExecutorService;
 
 import org.apache.camel.Exchange;
 import org.apache.camel.InvalidPayloadException;
@@ -131,33 +132,43 @@ public class FtpOperations implements 
RemoteFileOperations<FTPFile> {
                     client.getConnectTimeout());
         }
 
-        BlockingTask task = Tasks.foregroundTask()
-                .withBudget(Budgets.iterationBudget()
-                        
.withMaxIterations(Budgets.atLeastOnce(endpoint.getMaximumReconnectAttempts()))
-                        
.withInterval(Duration.ofMillis(endpoint.getReconnectDelay()))
-                        .build())
-                .build();
-
-        TaskPayload payload = new TaskPayload(configuration);
-
-        if (!task.run(endpoint.getCamelContext(), this::tryConnect, payload)) {
-            if (exchange != null) {
-                exchange.getIn().setHeader(FtpConstants.FTP_REPLY_CODE, 
client.getReplyCode());
-                exchange.getIn().setHeader(FtpConstants.FTP_REPLY_STRING, 
client.getReplyString());
-            }
+        ScheduledExecutorService ses = 
endpoint.getCamelContext().getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(this, "FtpReconnect");
+        try {
+            BlockingTask task = Tasks.backgroundTask()
+                    .withBudget(Budgets.iterationTimeBudget()
+                            
.withMaxIterations(Budgets.atLeastOnce(endpoint.getMaximumReconnectAttempts()))
+                            .withUnlimitedDuration()
+                            
.withInterval(Duration.ofMillis(endpoint.getReconnectDelay()))
+                            .build())
+                    .withScheduledExecutor(ses)
+                    .withName("FtpReconnect")
+                    .build();
+
+            TaskPayload payload = new TaskPayload(configuration);
+
+            if (!task.run(endpoint.getCamelContext(), this::tryConnect, 
payload)) {
+                if (exchange != null) {
+                    exchange.getIn().setHeader(FtpConstants.FTP_REPLY_CODE, 
client.getReplyCode());
+                    exchange.getIn().setHeader(FtpConstants.FTP_REPLY_STRING, 
client.getReplyString());
+                }
 
-            if (payload.exception != null) {
-                if (payload.exception instanceof 
GenericFileOperationFailedException genericFileOperationFailedException) {
-                    throw genericFileOperationFailedException;
+                if (payload.exception != null) {
+                    if (payload.exception instanceof 
GenericFileOperationFailedException genericFileOperationFailedException) {
+                        throw genericFileOperationFailedException;
+                    } else {
+                        throw new GenericFileOperationFailedException(
+                                client.getReplyCode(), 
client.getReplyString(), payload.exception.getMessage(),
+                                payload.exception);
+                    }
                 } else {
                     throw new GenericFileOperationFailedException(
-                            client.getReplyCode(), client.getReplyString(), 
payload.exception.getMessage(), payload.exception);
+                            client.getReplyCode(), client.getReplyString(),
+                            "Server refused connection");
                 }
-            } else {
-                throw new GenericFileOperationFailedException(
-                        client.getReplyCode(), client.getReplyString(),
-                        "Server refused connection");
             }
+        } finally {
+            
endpoint.getCamelContext().getExecutorServiceManager().shutdown(ses);
         }
 
         // we are now connected
diff --git 
a/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/SftpOperations.java
 
b/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/SftpOperations.java
index a7a09cc921ea..e921423868fa 100644
--- 
a/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/SftpOperations.java
+++ 
b/components/camel-ftp/src/main/java/org/apache/camel/component/file/remote/SftpOperations.java
@@ -35,6 +35,7 @@ import java.time.Duration;
 import java.util.Base64;
 import java.util.List;
 import java.util.Vector;
+import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.locks.Lock;
 import java.util.concurrent.locks.ReentrantLock;
 import java.util.regex.Pattern;
@@ -132,19 +133,28 @@ public class SftpOperations implements 
RemoteFileOperations<SftpRemoteFile> {
                 return true;
             }
 
-            BlockingTask task = Tasks.foregroundTask()
-                    .withBudget(Budgets.iterationBudget()
-                            
.withMaxIterations(Budgets.atLeastOnce(endpoint.getMaximumReconnectAttempts()))
-                            
.withInterval(Duration.ofMillis(endpoint.getReconnectDelay()))
-                            .build())
-                    .build();
-
-            TaskPayload payload = new TaskPayload(configuration);
-
-            if (!task.run(endpoint.getCamelContext(), this::tryConnect, 
payload)) {
-                throw new GenericFileOperationFailedException(
-                        "Cannot connect to " + 
configuration.remoteServerInformation(),
-                        payload.exception);
+            ScheduledExecutorService ses = 
endpoint.getCamelContext().getExecutorServiceManager()
+                    .newSingleThreadScheduledExecutor(this, "SftpReconnect");
+            try {
+                BlockingTask task = Tasks.backgroundTask()
+                        .withBudget(Budgets.iterationTimeBudget()
+                                
.withMaxIterations(Budgets.atLeastOnce(endpoint.getMaximumReconnectAttempts()))
+                                .withUnlimitedDuration()
+                                
.withInterval(Duration.ofMillis(endpoint.getReconnectDelay()))
+                                .build())
+                        .withScheduledExecutor(ses)
+                        .withName("SftpReconnect")
+                        .build();
+
+                TaskPayload payload = new TaskPayload(configuration);
+
+                if (!task.run(endpoint.getCamelContext(), this::tryConnect, 
payload)) {
+                    throw new GenericFileOperationFailedException(
+                            "Cannot connect to " + 
configuration.remoteServerInformation(),
+                            payload.exception);
+                }
+            } finally {
+                
endpoint.getCamelContext().getExecutorServiceManager().shutdown(ses);
             }
 
             configureBulkRequests();
diff --git 
a/components/camel-google/camel-google-pubsub/src/main/java/org/apache/camel/component/google/pubsub/GooglePubsubConsumer.java
 
b/components/camel-google/camel-google-pubsub/src/main/java/org/apache/camel/component/google/pubsub/GooglePubsubConsumer.java
index 49cc2ebccf2c..0fbf2a72f1d8 100644
--- 
a/components/camel-google/camel-google-pubsub/src/main/java/org/apache/camel/component/google/pubsub/GooglePubsubConsumer.java
+++ 
b/components/camel-google/camel-google-pubsub/src/main/java/org/apache/camel/component/google/pubsub/GooglePubsubConsumer.java
@@ -25,6 +25,7 @@ import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.CopyOnWriteArrayList;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.atomic.AtomicInteger;
 
 import com.google.api.core.AbstractApiService;
@@ -67,6 +68,7 @@ public class GooglePubsubConsumer extends DefaultConsumer 
implements ShutdownAwa
     private final Processor processor;
     private final AtomicInteger pendingExchanges = new AtomicInteger();
     private ExecutorService executor;
+    private ScheduledExecutorService taskExecutor;
     private final List<Subscriber> subscribers;
     private final Set<ApiFuture<PullResponse>> pendingSynchronousPullResponses;
     private final HeaderFilterStrategy headerFilterStrategy;
@@ -101,6 +103,8 @@ public class GooglePubsubConsumer extends DefaultConsumer 
implements ShutdownAwa
         }
 
         executor = endpoint.createExecutor(this);
+        taskExecutor = endpoint.getCamelContext().getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(this, "PubSubReconnectTask");
         for (int i = 0; i < endpoint.getConcurrentConsumers(); i++) {
             executor.submit(new SubscriberWrapper());
         }
@@ -122,6 +126,10 @@ public class GooglePubsubConsumer extends DefaultConsumer 
implements ShutdownAwa
 
     @Override
     protected void doShutdown() throws Exception {
+        if (taskExecutor != null) {
+            
getEndpoint().getCamelContext().getExecutorServiceManager().shutdown(taskExecutor);
+            taskExecutor = null;
+        }
         if (executor != null) {
             if (getEndpoint() != null && getEndpoint().getCamelContext() != 
null) {
                 
getEndpoint().getCamelContext().getExecutorServiceManager().shutdownGraceful(executor);
@@ -285,12 +293,14 @@ public class GooglePubsubConsumer extends DefaultConsumer 
implements ShutdownAwa
 
                     // Add backoff delay for recoverable errors to prevent 
tight loop
                     // We use initialDelay for the actual delay, and 
maxIterations(1) to run once
-                    Tasks.foregroundTask()
-                            .withBudget(Budgets.iterationBudget()
+                    Tasks.backgroundTask()
+                            .withBudget(Budgets.iterationTimeBudget()
                                     .withMaxIterations(1)
                                     .withInitialDelay(Duration.ofSeconds(5))
                                     .withInterval(Duration.ZERO)
+                                    .withUnlimitedDuration()
                                     .build())
+                            .withScheduledExecutor(taskExecutor)
                             .withName("PubSubRetryDelay")
                             .build()
                             .run(getEndpoint().getCamelContext(), () -> true);
diff --git 
a/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/seda/HazelcastSedaConsumer.java
 
b/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/seda/HazelcastSedaConsumer.java
index 81110a3ff295..da14472bfd15 100644
--- 
a/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/seda/HazelcastSedaConsumer.java
+++ 
b/components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/seda/HazelcastSedaConsumer.java
@@ -18,6 +18,7 @@ package org.apache.camel.component.hazelcast.seda;
 
 import java.time.Duration;
 import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 
 import com.hazelcast.collection.BaseQueue;
@@ -46,6 +47,7 @@ public class HazelcastSedaConsumer extends DefaultConsumer 
implements Runnable {
     private final HazelcastSedaEndpoint endpoint;
     private final AsyncProcessor processor;
     private ExecutorService executor;
+    private ScheduledExecutorService taskExecutor;
 
     public HazelcastSedaConsumer(final Endpoint endpoint, final Processor 
processor) {
         super(endpoint, processor);
@@ -55,6 +57,8 @@ public class HazelcastSedaConsumer extends DefaultConsumer 
implements Runnable {
 
     @Override
     protected void doStart() throws Exception {
+        taskExecutor = endpoint.getCamelContext().getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(this, 
"HazelcastSedaRecoveryTask");
         int concurrentConsumers = 
endpoint.getConfiguration().getConcurrentConsumers();
         executor = 
endpoint.getCamelContext().getExecutorServiceManager().newFixedThreadPool(this, 
endpoint.getEndpointUri(),
                 concurrentConsumers);
@@ -67,6 +71,10 @@ public class HazelcastSedaConsumer extends DefaultConsumer 
implements Runnable {
 
     @Override
     protected void doStop() throws Exception {
+        if (taskExecutor != null) {
+            
endpoint.getCamelContext().getExecutorServiceManager().shutdown(taskExecutor);
+            taskExecutor = null;
+        }
         if (executor != null) {
             
endpoint.getCamelContext().getExecutorServiceManager().shutdown(executor);
             executor = null;
@@ -153,12 +161,14 @@ public class HazelcastSedaConsumer extends 
DefaultConsumer implements Runnable {
                 getExceptionHandler().handleException("Error processing 
exchange", exchange, e);
                 // Use Camel's task API for error recovery delay instead of 
Thread.sleep()
                 // We use initialDelay for the actual delay, and 
maxIterations(1) to run once
-                Tasks.foregroundTask()
-                        .withBudget(Budgets.iterationBudget()
+                Tasks.backgroundTask()
+                        .withBudget(Budgets.iterationTimeBudget()
                                 .withMaxIterations(1)
                                 
.withInitialDelay(Duration.ofMillis(endpoint.getConfiguration().getOnErrorDelay()))
                                 .withInterval(Duration.ZERO)
+                                .withUnlimitedDuration()
                                 .build())
+                        .withScheduledExecutor(taskExecutor)
                         .withName("HazelcastSedaErrorRecoveryDelay")
                         .build()
                         .run(getEndpoint().getCamelContext(), () -> true);
diff --git 
a/components/camel-infinispan/camel-infinispan/src/main/java/org/apache/camel/component/infinispan/remote/InfinispanRemoteManager.java
 
b/components/camel-infinispan/camel-infinispan/src/main/java/org/apache/camel/component/infinispan/remote/InfinispanRemoteManager.java
index 68f7acdd19e9..092ce4e7f6f3 100644
--- 
a/components/camel-infinispan/camel-infinispan/src/main/java/org/apache/camel/component/infinispan/remote/InfinispanRemoteManager.java
+++ 
b/components/camel-infinispan/camel-infinispan/src/main/java/org/apache/camel/component/infinispan/remote/InfinispanRemoteManager.java
@@ -19,13 +19,14 @@ package org.apache.camel.component.infinispan.remote;
 import java.time.Duration;
 import java.util.Properties;
 import java.util.Set;
+import java.util.concurrent.ScheduledExecutorService;
 
 import org.apache.camel.CamelContext;
 import org.apache.camel.component.infinispan.InfinispanManager;
 import org.apache.camel.component.infinispan.InfinispanUtil;
 import 
org.apache.camel.component.infinispan.remote.embeddingstore.EmbeddingStoreUtil;
 import org.apache.camel.support.service.ServiceSupport;
-import org.apache.camel.support.task.ForegroundTask;
+import org.apache.camel.support.task.BackgroundTask;
 import org.apache.camel.support.task.Tasks;
 import org.apache.camel.support.task.budget.Budgets;
 import org.apache.camel.util.ObjectHelper;
@@ -148,38 +149,45 @@ public class InfinispanRemoteManager extends 
ServiceSupport implements Infinispa
      */
     private void registerSchemaWithRetry() throws Exception {
         Duration timeout = 
configuration.getEmbeddingStoreSchemaRegistrationTimeout();
-        ForegroundTask task = Tasks.foregroundTask()
-                .withName("infinispan-schema-registration")
-                .withBudget(Budgets.iterationTimeBudget()
-                        .withInterval(Duration.ofSeconds(1))
-                        .withMaxDuration(timeout)
-                        .build())
-                .build();
-
-        final boolean[] firstAttempt = { true };
-        boolean registered = task.run(camelContext, () -> {
-            try {
-                EmbeddingStoreUtil.registerSchema(configuration, 
cacheContainer);
-                return true;
-            } catch (HotRodClientException e) {
-                if (!isIllegalLifecycleStateException(e)) {
-                    throw e;
-                }
-                if (firstAttempt[0]) {
-                    firstAttempt[0] = false;
-                    LOG.info("Infinispan server not ready for schema 
registration, will retry for up to {}: {}",
-                            timeout, e.getMessage());
-                } else {
-                    LOG.debug("Schema registration failed (server not ready), 
retrying: {}", e.getMessage());
+        ScheduledExecutorService ses = camelContext.getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(this, 
"infinispan-schema-registration-task");
+        try {
+            BackgroundTask task = Tasks.backgroundTask()
+                    .withBudget(Budgets.iterationTimeBudget()
+                            .withInterval(Duration.ofSeconds(1))
+                            .withMaxDuration(timeout)
+                            .build())
+                    .withScheduledExecutor(ses)
+                    .withName("infinispan-schema-registration")
+                    .build();
+
+            final boolean[] firstAttempt = { true };
+            boolean registered = task.run(camelContext, () -> {
+                try {
+                    EmbeddingStoreUtil.registerSchema(configuration, 
cacheContainer);
+                    return true;
+                } catch (HotRodClientException e) {
+                    if (!isIllegalLifecycleStateException(e)) {
+                        throw e;
+                    }
+                    if (firstAttempt[0]) {
+                        firstAttempt[0] = false;
+                        LOG.info("Infinispan server not ready for schema 
registration, will retry for up to {}: {}",
+                                timeout, e.getMessage());
+                    } else {
+                        LOG.debug("Schema registration failed (server not 
ready), retrying: {}", e.getMessage());
+                    }
+                    return false;
                 }
-                return false;
-            }
-        });
+            });
 
-        if (!registered) {
-            throw new IllegalStateException(
-                    "Failed to register Infinispan schema after " + timeout
-                                            + " of retries. The server may not 
be fully started.");
+            if (!registered) {
+                throw new IllegalStateException(
+                        "Failed to register Infinispan schema after " + timeout
+                                                + " of retries. The server may 
not be fully started.");
+            }
+        } finally {
+            camelContext.getExecutorServiceManager().shutdown(ses);
         }
     }
 
diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaFetchRecords.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaFetchRecords.java
index 493fa245fa71..68ebb7d8aba7 100644
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaFetchRecords.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaFetchRecords.java
@@ -19,6 +19,7 @@ package org.apache.camel.component.kafka;
 import java.time.Duration;
 import java.util.Collection;
 import java.util.Properties;
+import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.concurrent.locks.ReentrantLock;
@@ -43,7 +44,7 @@ import 
org.apache.camel.component.kafka.consumer.support.subcription.DefaultSubs
 import 
org.apache.camel.component.kafka.consumer.support.subcription.SubscribeAdapter;
 import org.apache.camel.component.kafka.consumer.support.subcription.TopicInfo;
 import org.apache.camel.support.BridgeExceptionHandlerToErrorHandler;
-import org.apache.camel.support.task.ForegroundTask;
+import org.apache.camel.support.task.BackgroundTask;
 import org.apache.camel.support.task.TaskRunFailureException;
 import org.apache.camel.support.task.Tasks;
 import org.apache.camel.support.task.budget.Budgets;
@@ -129,6 +130,10 @@ public class KafkaFetchRecords implements Runnable {
             return;
         }
 
+        ScheduledExecutorService taskExecutor = 
kafkaConsumer.getEndpoint().getCamelContext()
+                .getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(kafkaConsumer, 
"KafkaFetchRecordsTask-" + threadId);
+
         do {
             terminated = false;
 
@@ -141,14 +146,16 @@ public class KafkaFetchRecords implements Runnable {
 
                 // task that deals with creating kafka consumer
                 currentBackoffInterval = 
kafkaConsumer.getEndpoint().getComponent().getCreateConsumerBackoffInterval();
-                ForegroundTask task = Tasks.foregroundTask()
-                        .withName("Create KafkaConsumer")
-                        .withBudget(Budgets.iterationBudget()
+                BackgroundTask task = Tasks.backgroundTask()
+                        .withBudget(Budgets.iterationTimeBudget()
                                 .withMaxIterations(
                                         
kafkaConsumer.getEndpoint().getComponent().getCreateConsumerBackoffMaxAttempts())
                                 .withInitialDelay(Duration.ZERO)
                                 
.withInterval(Duration.ofMillis(currentBackoffInterval))
+                                .withUnlimitedDuration()
                                 .build())
+                        .withScheduledExecutor(taskExecutor)
+                        .withName("Create KafkaConsumer")
                         .build();
                 boolean success = 
task.run(kafkaConsumer.getEndpoint().getCamelContext(), 
this::createConsumerTask);
                 if (!success) {
@@ -161,14 +168,16 @@ public class KafkaFetchRecords implements Runnable {
 
                 // task that deals with subscribing kafka consumer
                 currentBackoffInterval = 
kafkaConsumer.getEndpoint().getComponent().getSubscribeConsumerBackoffInterval();
-                task = Tasks.foregroundTask()
-                        .withName("Subscribe KafkaConsumer")
-                        .withBudget(Budgets.iterationBudget()
+                task = Tasks.backgroundTask()
+                        .withBudget(Budgets.iterationTimeBudget()
                                 .withMaxIterations(
                                         
kafkaConsumer.getEndpoint().getComponent().getSubscribeConsumerBackoffMaxAttempts())
                                 .withInitialDelay(Duration.ZERO)
                                 
.withInterval(Duration.ofMillis(currentBackoffInterval))
+                                .withUnlimitedDuration()
                                 .build())
+                        .withScheduledExecutor(taskExecutor)
+                        .withName("Subscribe KafkaConsumer")
                         .build();
                 success = 
task.run(kafkaConsumer.getEndpoint().getCamelContext(), 
this::initializeConsumerTask);
                 if (!success) {
@@ -194,10 +203,11 @@ public class KafkaFetchRecords implements Runnable {
             LOG.info("Terminating KafkaConsumer thread {} receiving from {}", 
threadId, getPrintableTopic());
         }
 
+        
kafkaConsumer.getEndpoint().getCamelContext().getExecutorServiceManager().shutdown(taskExecutor);
         safeConsumerClose();
     }
 
-    private void setupInitializeErrorException(ForegroundTask task, int max) {
+    private void setupInitializeErrorException(BackgroundTask task, int max) {
         String time = TimeUtils.printDuration(task.elapsed(), true);
         String topic = getPrintableTopic();
         String msg = "Gave up subscribing 
org.apache.kafka.clients.consumer.KafkaConsumer " +
@@ -206,7 +216,7 @@ public class KafkaFetchRecords implements Runnable {
         setLastError(new KafkaConsumerFatalException(msg, lastError));
     }
 
-    private void setupCreateConsumerException(ForegroundTask task, int max) {
+    private void setupCreateConsumerException(BackgroundTask task, int max) {
         String time = TimeUtils.printDuration(task.elapsed(), true);
         String topic = getPrintableTopic();
         String msg = "Gave up creating 
org.apache.kafka.clients.consumer.KafkaConsumer "
diff --git 
a/components/camel-mllp/src/main/java/org/apache/camel/component/mllp/internal/TcpServerBindThread.java
 
b/components/camel-mllp/src/main/java/org/apache/camel/component/mllp/internal/TcpServerBindThread.java
index 0dc572f9f081..70a758a0e79a 100644
--- 
a/components/camel-mllp/src/main/java/org/apache/camel/component/mllp/internal/TcpServerBindThread.java
+++ 
b/components/camel-mllp/src/main/java/org/apache/camel/component/mllp/internal/TcpServerBindThread.java
@@ -22,6 +22,7 @@ import java.net.ServerSocket;
 import java.net.SocketException;
 import java.security.GeneralSecurityException;
 import java.time.Duration;
+import java.util.concurrent.ScheduledExecutorService;
 
 import javax.net.ssl.SSLContext;
 import javax.net.ssl.SSLServerSocketFactory;
@@ -100,19 +101,27 @@ public class TcpServerBindThread extends Thread {
     }
 
     private void doAccept(ServerSocket serverSocket, InetSocketAddress 
socketAddress) {
-        BlockingTask task = Tasks.foregroundTask()
-                .withBudget(Budgets.iterationTimeBudget()
-                        
.withMaxDuration(Duration.ofMillis(consumer.getConfiguration().getBindTimeout()))
-                        
.withInterval(Duration.ofMillis(consumer.getConfiguration().getBindRetryInterval()))
-                        .build())
-                .withName("mllp-tcp-server-accept")
-                .build();
-
-        if (task.run(consumer.getEndpoint().getCamelContext(), () -> 
doBind(serverSocket, socketAddress))) {
-            consumer.startAcceptThread(serverSocket);
-        } else {
-            log.error("Failed to bind to address {} within timeout {}", 
socketAddress,
-                    consumer.getConfiguration().getBindTimeout());
+        ScheduledExecutorService ses = consumer.getEndpoint().getCamelContext()
+                .getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(consumer, 
"mllp-tcp-server-bind-task");
+        try {
+            BlockingTask task = Tasks.backgroundTask()
+                    .withBudget(Budgets.iterationTimeBudget()
+                            
.withMaxDuration(Duration.ofMillis(consumer.getConfiguration().getBindTimeout()))
+                            
.withInterval(Duration.ofMillis(consumer.getConfiguration().getBindRetryInterval()))
+                            .build())
+                    .withScheduledExecutor(ses)
+                    .withName("mllp-tcp-server-accept")
+                    .build();
+
+            if (task.run(consumer.getEndpoint().getCamelContext(), () -> 
doBind(serverSocket, socketAddress))) {
+                consumer.startAcceptThread(serverSocket);
+            } else {
+                log.error("Failed to bind to address {} within timeout {}", 
socketAddress,
+                        consumer.getConfiguration().getBindTimeout());
+            }
+        } finally {
+            
consumer.getEndpoint().getCamelContext().getExecutorServiceManager().shutdown(ses);
         }
     }
 
diff --git 
a/components/camel-mongodb-gridfs/src/main/java/org/apache/camel/component/mongodb/gridfs/GridFsConsumer.java
 
b/components/camel-mongodb-gridfs/src/main/java/org/apache/camel/component/mongodb/gridfs/GridFsConsumer.java
index 3ca6c9fe2b52..c8d40b2faad7 100644
--- 
a/components/camel-mongodb-gridfs/src/main/java/org/apache/camel/component/mongodb/gridfs/GridFsConsumer.java
+++ 
b/components/camel-mongodb-gridfs/src/main/java/org/apache/camel/component/mongodb/gridfs/GridFsConsumer.java
@@ -20,6 +20,7 @@ import java.io.InputStream;
 import java.time.Duration;
 import java.util.Date;
 import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
 
 import com.mongodb.BasicDBObject;
 import com.mongodb.client.MongoCollection;
@@ -36,7 +37,6 @@ import org.apache.camel.support.DefaultConsumer;
 import org.apache.camel.support.task.BlockingTask;
 import org.apache.camel.support.task.Tasks;
 import org.apache.camel.support.task.budget.Budgets;
-import org.apache.camel.support.task.budget.IterationBoundedBudget;
 import org.bson.Document;
 import org.bson.conversions.Bson;
 
@@ -49,6 +49,7 @@ import static 
org.apache.camel.component.mongodb.gridfs.GridFsConstants.PERSISTE
 
 public class GridFsConsumer extends DefaultConsumer implements Runnable {
     private final GridFsEndpoint endpoint;
+    private volatile ScheduledExecutorService taskExecutor;
     private volatile ExecutorService executor;
 
     public GridFsConsumer(GridFsEndpoint endpoint, Processor processor) {
@@ -59,6 +60,10 @@ public class GridFsConsumer extends DefaultConsumer 
implements Runnable {
     @Override
     protected void doStop() throws Exception {
         super.doStop();
+        if (taskExecutor != null) {
+            
endpoint.getCamelContext().getExecutorServiceManager().shutdown(taskExecutor);
+            taskExecutor = null;
+        }
         if (executor != null) {
             
endpoint.getCamelContext().getExecutorServiceManager().shutdown(executor);
             executor = null;
@@ -68,6 +73,8 @@ public class GridFsConsumer extends DefaultConsumer 
implements Runnable {
     @Override
     protected void doStart() throws Exception {
         super.doStart();
+        taskExecutor = endpoint.getCamelContext().getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(this, "GridFsPollingTask");
         executor = 
endpoint.getCamelContext().getExecutorServiceManager().newFixedThreadPool(this, 
endpoint.getEndpointUri(),
                 1);
         executor.execute(this);
@@ -107,12 +114,15 @@ public class GridFsConsumer extends DefaultConsumer 
implements Runnable {
             fromDate = new Date();
         }
 
-        BlockingTask task = Tasks.foregroundTask()
-                .withBudget(Budgets.iterationBudget()
-                        
.withMaxIterations(IterationBoundedBudget.UNLIMITED_ITERATIONS)
+        BlockingTask task = Tasks.backgroundTask()
+                .withBudget(Budgets.iterationTimeBudget()
+                        .withMaxIterations(Integer.MAX_VALUE)
                         .withInterval(Duration.ofMillis(endpoint.getDelay()))
                         
.withInitialDelay(Duration.ofMillis(endpoint.getInitialDelay()))
+                        .withUnlimitedDuration()
                         .build())
+                .withScheduledExecutor(taskExecutor)
+                .withName("GridFsPolling")
                 .build();
 
         MongoCollection<Document> finalPtsCollection = ptsCollection;
diff --git 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/streaming/SubscriptionHelper.java
 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/streaming/SubscriptionHelper.java
index 9a6d4b13e209..1bf08486f075 100644
--- 
a/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/streaming/SubscriptionHelper.java
+++ 
b/components/camel-salesforce/camel-salesforce-component/src/main/java/org/apache/camel/component/salesforce/internal/streaming/SubscriptionHelper.java
@@ -30,6 +30,7 @@ import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.locks.Lock;
@@ -93,6 +94,7 @@ public class SubscriptionHelper extends ServiceSupport {
 
     BayeuxClient client;
 
+    private ScheduledExecutorService taskExecutor;
     private final SalesforceComponent component;
     private SalesforceSession session;
 
@@ -151,12 +153,14 @@ public class SubscriptionHelper extends ServiceSupport {
                 } else {
                     LOG.debug("Pausing for {} msecs before handshake retry", 
backoff);
                     if (backoff > 0) {
-                        Tasks.foregroundTask()
-                                .withBudget(Budgets.iterationBudget()
+                        Tasks.backgroundTask()
+                                .withBudget(Budgets.iterationTimeBudget()
                                         .withMaxIterations(1)
                                         
.withInitialDelay(Duration.ofMillis(backoff))
                                         .withInterval(Duration.ZERO)
+                                        .withUnlimitedDuration()
                                         .build())
+                                .withScheduledExecutor(taskExecutor)
                                 .withName("SalesforceHandshakeRetryDelay")
                                 .build()
                                 .run(component.getCamelContext(), () -> true);
@@ -286,12 +290,14 @@ public class SubscriptionHelper extends ServiceSupport {
                 LOG.debug("Pausing for {} msecs before subscribe attempt", 
backoff);
                 // Use Camel's task API for backoff delay instead of 
Thread.sleep()
                 // We use initialDelay for the actual delay, and 
maxIterations(1) to run once
-                Tasks.foregroundTask()
-                        .withBudget(Budgets.iterationBudget()
+                Tasks.backgroundTask()
+                        .withBudget(Budgets.iterationTimeBudget()
                                 .withMaxIterations(1)
                                 .withInitialDelay(Duration.ofMillis(backoff))
                                 .withInterval(Duration.ZERO)
+                                .withUnlimitedDuration()
                                 .build())
+                        .withScheduledExecutor(taskExecutor)
                         .withName("SalesforceSubscribeRetryDelay")
                         .build()
                         .run(component.getCamelContext(), () -> true);
@@ -330,6 +336,9 @@ public class SubscriptionHelper extends ServiceSupport {
             throw new CamelException("Lazy login is not supported by 
salesforce consumers.");
         }
 
+        taskExecutor = component.getCamelContext().getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(this, 
"SalesforceReconnectTask");
+
         // create CometD client
         client = createClient(component, session);
 
@@ -397,6 +406,11 @@ public class SubscriptionHelper extends ServiceSupport {
 
     @Override
     protected void doStop() throws Exception {
+        if (taskExecutor != null) {
+            
component.getCamelContext().getExecutorServiceManager().shutdown(taskExecutor);
+            taskExecutor = null;
+        }
+
         closeChannel(META_CONNECT);
         closeChannel(META_SUBSCRIBE);
         closeChannel(META_HANDSHAKE);
diff --git 
a/components/camel-zookeeper/src/main/java/org/apache/camel/component/zookeeper/ZooKeeperConsumer.java
 
b/components/camel-zookeeper/src/main/java/org/apache/camel/component/zookeeper/ZooKeeperConsumer.java
index 16f81de19fff..2ddf50e272f4 100644
--- 
a/components/camel-zookeeper/src/main/java/org/apache/camel/component/zookeeper/ZooKeeperConsumer.java
+++ 
b/components/camel-zookeeper/src/main/java/org/apache/camel/component/zookeeper/ZooKeeperConsumer.java
@@ -19,6 +19,7 @@ package org.apache.camel.component.zookeeper;
 import java.time.Duration;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ScheduledExecutorService;
 
 import org.apache.camel.Exchange;
 import org.apache.camel.Processor;
@@ -53,6 +54,7 @@ public class ZooKeeperConsumer extends DefaultConsumer {
     private ZooKeeperConfiguration configuration;
     private LinkedBlockingQueue<ZooKeeperOperation> operations = new 
LinkedBlockingQueue<>();
     private ExecutorService executor;
+    private ScheduledExecutorService taskExecutor;
     private volatile boolean shuttingDown;
 
     public ZooKeeperConsumer(ZooKeeperEndpoint endpoint, Processor processor) {
@@ -72,6 +74,8 @@ public class ZooKeeperConsumer extends DefaultConsumer {
         initializeConsumer();
         executor = 
getEndpoint().getCamelContext().getExecutorServiceManager().newFixedThreadPool(this,
                 "Camel-Zookeeper OperationsExecutor", 1);
+        taskExecutor = 
getEndpoint().getCamelContext().getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(this, 
"ZooKeeperReconnectTask");
 
         OperationsExecutor opsService = new OperationsExecutor();
         executor.submit(opsService);
@@ -84,6 +88,10 @@ public class ZooKeeperConsumer extends DefaultConsumer {
         if (LOG.isTraceEnabled()) {
             LOG.trace(String.format("Shutting down zookeeper consumer of 
'%s'", configuration.getPath()));
         }
+        if (taskExecutor != null) {
+            
getEndpoint().getCamelContext().getExecutorServiceManager().shutdown(taskExecutor);
+            taskExecutor = null;
+        }
         
getEndpoint().getCamelContext().getExecutorServiceManager().shutdown(executor);
         zkm.shutdown();
     }
@@ -179,12 +187,14 @@ public class ZooKeeperConsumer extends DefaultConsumer {
                 if (isRunAllowed()) {
                     // Use Camel's task API for reconnection backoff delay 
instead of Thread.sleep()
                     // We use initialDelay for the actual delay, and 
maxIterations(1) to run once
-                    Tasks.foregroundTask()
-                            .withBudget(Budgets.iterationBudget()
+                    Tasks.backgroundTask()
+                            .withBudget(Budgets.iterationTimeBudget()
                                     .withMaxIterations(1)
                                     
.withInitialDelay(Duration.ofMillis(configuration.getBackoff()))
                                     .withInterval(Duration.ZERO)
+                                    .withUnlimitedDuration()
                                     .build())
+                            .withScheduledExecutor(taskExecutor)
                             .withName("ZooKeeperReconnectBackoff")
                             .build()
                             .run(getEndpoint().getCamelContext(), () -> true);


Reply via email to