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 4c51b5f..f6b64ed 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 @@ -825,9 +825,6 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i case EXPIRED: this.task = () -> { try { - if (logger.isInfoEnabled() && this.state == ConsumerState.NEW) { - logger.info("The [" + this + "] has been started."); - } if (this.shardOffset.isReset()) { this.checkpointer.remove(); } @@ -838,7 +835,9 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i this.shardOffset.setIteratorType(ShardIteratorType.AFTER_SEQUENCE_NUMBER); } } - + if (logger.isInfoEnabled() && this.state == ConsumerState.NEW) { + logger.info("The [" + this + "] has been started."); + } GetShardIteratorRequest shardIteratorRequest = this.shardOffset.toShardIteratorRequest(); this.shardIterator = KinesisMessageDrivenChannelAdapter.this.amazonKinesis