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(); } } }
