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
This commit is contained in:
@@ -77,6 +77,8 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport {
|
||||
|
||||
private static final ThreadLocal<AttributeAccessor> 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<byte[]> 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();
|
||||
|
||||
Reference in New Issue
Block a user