From 2eb0ba2c3ffdc6b948172b02a5da2b55b4857199 Mon Sep 17 00:00:00 2001 From: abilan Date: Wed, 30 Nov 2022 11:28:46 -0500 Subject: [PATCH] GH-210: Short-circuit Kinesis consumer for stop Fixes https://github.com/spring-projects/spring-integration-aws/issues/210 To avoid extra cycles for tasks and locks renewal check for a closed shard just after `getShardIterator()` request in a `NEW` consumer task **Cherry-pick to `2.5.x`** --- .../inbound/kinesis/KinesisMessageDrivenChannelAdapter.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java index ad7475b..34d6c35 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java @@ -974,6 +974,10 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport .amazonKinesis .getShardIterator(shardIteratorRequest) .getShardIterator(); + if (this.shardIterator == null) { + // The shard is closed - stop consumer + this.state = ConsumerState.STOP; + } if (ConsumerState.STOP != this.state) { this.state = ConsumerState.CONSUME; }