From 8e7e1a02b6c64544554e6a367fa7ab1c931c0d63 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 4 Apr 2017 19:54:57 -0400 Subject: [PATCH] GH-63: Kinesis: fix logic to skip closed shards Fixes spring-projects/spring-integration-aws#63 The `CLOSED` shards have `endingSequenceNumber` value and they can't be considered for consuming independently of the value in the `checkpoint` * Introduce local `shardsToConsume` variable in the `KinesisMessageDrivenChannelAdapter#populateShardsForStream()` to store only `OPEN` shards for future consumption * Fix the logic to determine the last `shardId` for a subsequent `describeStream()` request based on the shards result exactly from the `describeStreamResult.getStreamDescription().getShards()`, not already filtered `shardsToConsume` --- .../KinesisMessageDrivenChannelAdapter.java | 32 ++++++------------- 1 file changed, 10 insertions(+), 22 deletions(-) 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 90ccaa2..04537e1 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 @@ -16,7 +16,6 @@ package org.springframework.integration.aws.inbound.kinesis; -import java.math.BigInteger; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; @@ -427,7 +426,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i public void run() { try { int describeStreamRetries = 0; - List shards = new ArrayList<>(); + List shardsToConsume = new ArrayList<>(); String exclusiveStartShardId = null; while (true) { @@ -447,8 +446,9 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i KinesisMessageDrivenChannelAdapter.this.describeStreamBackoff + "] millis."); } - if (describeStreamResult == null || !StreamStatus.ACTIVE.toString() - .equals(describeStreamResult.getStreamDescription().getStreamStatus())) { + if (describeStreamResult == null || + !StreamStatus.ACTIVE.toString() + .equals(describeStreamResult.getStreamDescription().getStreamStatus())) { if (describeStreamRetries++ > KinesisMessageDrivenChannelAdapter.this.describeStreamRetries) { ResourceNotFoundException resourceNotFoundException = @@ -468,24 +468,12 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i } } - for (Shard shard : describeStreamResult.getStreamDescription().getShards()) { - String endingSequenceNumber = shard.getSequenceNumberRange().getEndingSequenceNumber(); - if (endingSequenceNumber != null) { - String key = KinesisMessageDrivenChannelAdapter.this.consumerGroup + - ":" + stream + - ":" + shard.getShardId(); - String checkpoint = - KinesisMessageDrivenChannelAdapter.this.checkpointStore.get(key); - - if (checkpoint != null && - new BigInteger(endingSequenceNumber) - .compareTo(new BigInteger(checkpoint)) <= 0) { - // Skip CLOSED shard which has been read before according a checkpoint - continue; - } + List shards = describeStreamResult.getStreamDescription().getShards(); + for (Shard shard : shards) { + // Check if the shard is still open. Open shards do not have an ending sequence number. + if (shard.getSequenceNumberRange().getEndingSequenceNumber() == null) { + shardsToConsume.add(shard); } - - shards.add(shard); } if (describeStreamResult.getStreamDescription().getHasMoreShards()) { @@ -496,7 +484,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i } } - for (Shard shard : shards) { + for (Shard shard : shardsToConsume) { KinesisShardOffset shardOffset = new KinesisShardOffset(KinesisMessageDrivenChannelAdapter.this.streamInitialSequence); shardOffset.setShard(shard.getShardId());