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

davsclaus pushed a commit to branch fix/CAMEL-24156
in repository https://gitbox.apache.org/repos/asf/camel.git

commit f319494b24b9742f652f67e6726bfe28ac12b5f0
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 17:21:15 2026 +0200

    CAMEL-24156: camel-aws2-kinesis - Fix connection close reset and resume 
concurrency
    
    - Reset client fields to null in KinesisConnection.close() so subsequent
      getClient() calls create a fresh client instead of returning a closed one
    - Always create a new KinesisResumeAction per shard instead of reusing a
      shared registry instance that gets mutated concurrently from 
parallelStream
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../camel/component/aws2/kinesis/Kinesis2Consumer.java  | 10 +---------
 .../camel/component/aws2/kinesis/KinesisConnection.java | 17 ++++++++++++-----
 2 files changed, 13 insertions(+), 14 deletions(-)

diff --git 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Consumer.java
 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Consumer.java
index e452792ccbac..b3bc7e1f2172 100644
--- 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Consumer.java
+++ 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Consumer.java
@@ -362,15 +362,7 @@ public class Kinesis2Consumer extends 
ScheduledBatchPollingConsumer implements R
     }
 
     private KinesisResumeAction resolveResumeAction(String shardId, 
GetShardIteratorRequest.Builder req) {
-        KinesisResumeAction action
-                = 
getEndpoint().getCamelContext().getRegistry().lookupByNameAndType(Kinesis2Constants.RESUME_ACTION,
-                        KinesisResumeAction.class);
-        if (ObjectHelper.isEmpty(action)) {
-            action = new KinesisResumeAction(req);
-        } else {
-            action.setBuilder(req);
-        }
-
+        KinesisResumeAction action = new KinesisResumeAction(req);
         action.setShardId(shardId);
         action.setStreamName(getEndpoint().getConfiguration().getStreamName());
         return action;
diff --git 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/KinesisConnection.java
 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/KinesisConnection.java
index b9e888f3728e..64d32eaf34a1 100644
--- 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/KinesisConnection.java
+++ 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/KinesisConnection.java
@@ -73,11 +73,18 @@ public class KinesisConnection implements Closeable {
 
     @Override
     public void close() throws IOException {
-        if (ObjectHelper.isNotEmpty(kinesisClient)) {
-            kinesisClient.close();
-        }
-        if (ObjectHelper.isNotEmpty(kinesisAsyncClient)) {
-            kinesisAsyncClient.close();
+        lock.lock();
+        try {
+            if (ObjectHelper.isNotEmpty(kinesisClient)) {
+                kinesisClient.close();
+                kinesisClient = null;
+            }
+            if (ObjectHelper.isNotEmpty(kinesisAsyncClient)) {
+                kinesisAsyncClient.close();
+                kinesisAsyncClient = null;
+            }
+        } finally {
+            lock.unlock();
         }
     }
 }

Reply via email to