From 5e9f76a1664ec338c3f0109ce9960aca6cdd3095 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 10 Sep 2019 13:36:24 -0400 Subject: [PATCH] GH-151: Recreate KCL Worker on each restart Fixes https://github.com/spring-projects/spring-integration-aws/issues/151 **Cherry-pick to 2.2.x** # Conflicts: # src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java --- .../KclMessageDrivenChannelAdapter.java | 40 ++++++++++++------- 1 file changed, 26 insertions(+), 14 deletions(-) diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java index c2fa167..6e82d1d 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java @@ -77,6 +77,8 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { private static final ThreadLocal attributesHolder = new ThreadLocal<>(); + private final RecordProcessorFactory recordProcessorFactory = new RecordProcessorFactory(); + private final String stream; private final AmazonKinesis kinesisClient; @@ -93,7 +95,7 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { private InboundMessageMapper embeddedHeadersMapper; - private Worker scheduler; + private KinesisClientLibConfiguration config; private InitialPositionInStream streamInitialSequence = InitialPositionInStream.LATEST; @@ -109,6 +111,8 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { private boolean bindSourceRecord; + private volatile Worker scheduler; + public KclMessageDrivenChannelAdapter(String streams) { this(streams, AmazonKinesisClientBuilder.defaultClient(), AmazonCloudWatchClientBuilder.defaultClient(), AmazonDynamoDBClientBuilder.defaultClient(), @@ -207,14 +211,14 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { protected void onInit() { super.onInit(); - KinesisClientLibConfiguration config = - new KinesisClientLibConfiguration( - this.consumerGroup, + this.config = + new KinesisClientLibConfiguration(this.consumerGroup, this.stream, null, this.streamInitialSequence, this.kinesisProxyCredentialsProvider, - null, null, + null, + null, KinesisClientLibConfiguration.DEFAULT_FAILOVER_TIME_MILLIS, this.workerId, KinesisClientLibConfiguration.DEFAULT_MAX_RECORDS, @@ -232,20 +236,22 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { KinesisClientLibConfiguration.DEFAULT_VALIDATE_SEQUENCE_NUMBER_BEFORE_CHECKPOINTING, null, KinesisClientLibConfiguration.DEFAULT_SHUTDOWN_GRACE_MILLIS); - - this.scheduler = new Worker.Builder() - .kinesisClient(this.kinesisClient) - .dynamoDBClient(this.dynamoDBClient) - .cloudWatchClient(this.cloudWatchClient) - .recordProcessorFactory(new RecordProcessorFactory()) - .execService(new ExecutorServiceAdapter(this.executor)) - .config(config) - .build(); } @Override protected void doStart() { super.doStart(); + this.scheduler = + new Worker + .Builder() + .kinesisClient(this.kinesisClient) + .dynamoDBClient(this.dynamoDBClient) + .cloudWatchClient(this.cloudWatchClient) + .recordProcessorFactory(this.recordProcessorFactory) + .execService(new ExecutorServiceAdapter(this.executor)) + .config(this.config) + .build(); + this.executor.execute(this.scheduler); } @@ -260,6 +266,12 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { } + @Override + public void destroy() { + super.destroy(); + this.scheduler.shutdown(); + } + @Override protected AttributeAccessor getErrorMessageAttributes(org.springframework.messaging.Message message) { AttributeAccessor attributes = attributesHolder.get();