From 2b25b7e5d75551ea507152131ca4c2b664298cff Mon Sep 17 00:00:00 2001 From: Greg Eales Date: Mon, 19 Oct 2020 08:37:11 -0400 Subject: [PATCH] GH-177: Checkpoint closed shards for Kinesis Fixes https://github.com/spring-projects/spring-integration-aws/issues/177 * When the end of a shard is detected, checkpoint with the `endingSequenceNumber` * Add a test * code review * Default to directly checkpointing closed shards except when in manual checkpoint mode * update tests * `endingSequenceNumber` should be higher than the last record's sequence number * Checkpoint when in manual mode if shard was empty --- .../KinesisMessageDrivenChannelAdapter.java | 24 +++++-- ...nesisMessageDrivenChannelAdapterTests.java | 71 ++++++++++++++----- 2 files changed, 74 insertions(+), 21 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 f8a2502..43561a0 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 @@ -91,6 +91,8 @@ import com.amazonaws.services.kinesis.model.ShardIteratorType; * @author Krzysztof Witkowski * @author Hervé Fortin * @author Dirk Bonhomme + * @author Greg Eales + * * @since 1.1 */ @ManagedResource @@ -293,7 +295,6 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport * will be processed sequentially. In other words each shard is tied with the particular thread. * By default the concurrency is unlimited and shard is processed in the {@link #consumerExecutor} * directly. - * * @param concurrency the concurrency maximum number */ public void setConcurrency(int concurrency) { @@ -303,7 +304,6 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport /** * The sleep interval in milliseconds used in the main loop between shards polling cycles. * Defaults to {@code 1000}l minimum {@code 250}. - * * @param idleBetweenPolls the interval to sleep between shards polling cycles. */ public void setIdleBetweenPolls(int idleBetweenPolls) { @@ -313,7 +313,6 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport /** * Specify an {@link InboundMessageMapper} to extract message headers embedded into the record * data. - * * @param embeddedHeadersMapper the {@link InboundMessageMapper} to use. * @since 2.0 */ @@ -324,7 +323,6 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport /** * Specify a {@link LockRegistry} for an exclusive access to provided streams. This is not used * when shards-based configuration is provided. - * * @param lockRegistry the {@link LockRegistry} to use. * @since 2.0 */ @@ -335,7 +333,6 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport /** * Set to true to bind the source consumer record in the header named {@link * IntegrationMessageHeaderAccessor#SOURCE_DATA}. Does not apply to batch listeners. - * * @param bindSourceRecord true to bind. * @since 2.2 */ @@ -347,6 +344,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport * Specify a {@link Function Function<List<Shard>, List<Shard>>} to filter the shards which will * be read from. * @param shardListFilter the filter {@link Function Function<List<Shard>, List<Shard>>} + * @since 2.3.4 */ public void setShardListFilter(Function, List> shardListFilter) { this.shardListFilter = shardListFilter; @@ -1027,6 +1025,22 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport .remove(this.key); } // Shard is closed: nothing to consume any more. + // Checkpoint endingSequenceNumber to ensure shard is marked exhausted. + // If in CheckpointMode.manual, only checkpoint if lastCheckpointValue is also null, as this + // means that no records have ever been read and so the shard was empty + if (!CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode) + || this.checkpointer.getLastCheckpointValue() == null) { + for (Shard shard : readShardList(this.shardOffset.getStream())) { + if (shard.getShardId().equals(this.shardOffset.getShard())) { + String endingSequenceNumber = + shard.getSequenceNumberRange().getEndingSequenceNumber(); + if (endingSequenceNumber != null) { + this.checkpointer.checkpoint(endingSequenceNumber); + } + break; + } + } + } // Resharding is possible. if (KinesisMessageDrivenChannelAdapter.this.applicationEventPublisher != null) { KinesisMessageDrivenChannelAdapter.this.applicationEventPublisher.publishEvent( diff --git a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java index 7b4a1b3..6066c83 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java @@ -74,6 +74,8 @@ import com.amazonaws.services.kinesis.model.Shard; /** * @author Artem Bilan * @author Matthias Wesolowski + * @author Greg Eales + * * @since 1.1 */ @SpringJUnitConfig @@ -93,6 +95,9 @@ public class KinesisMessageDrivenChannelAdapterTests { @Autowired private MetadataStore checkpointStore; + @Autowired + private MetadataStore reshardingCheckpointStore; + @Autowired private KinesisMessageDrivenChannelAdapter reshardingChannelAdapter; @@ -108,7 +113,7 @@ public class KinesisMessageDrivenChannelAdapterTests { } @Test - @SuppressWarnings({"unchecked", "rawtypes"}) + @SuppressWarnings({ "unchecked", "rawtypes" }) void testKinesisMessageDrivenChannelAdapter() { this.kinesisMessageDrivenChannelAdapter.start(); final Set shardOffsets = TestUtils.getPropertyValue(this.kinesisMessageDrivenChannelAdapter, @@ -219,11 +224,14 @@ public class KinesisMessageDrivenChannelAdapterTests { this.reshardingChannelAdapter.stop(); + assertThat(this.reshardingCheckpointStore.get("SpringIntegration:streamForResharding:closedEmptyShard5")) + .isEqualTo("50"); + KinesisShardEndedEvent kinesisShardEndedEvent = this.config.shardEndedEventReference.get(); assertThat(kinesisShardEndedEvent).isNotNull() .extracting(KinesisShardEndedEvent::getShardKey) - .isEqualTo("SpringIntegration:streamForResharding:closedShard4"); + .isEqualTo("SpringIntegration:streamForResharding:closedEmptyShard5"); } @Configuration @@ -336,39 +344,51 @@ public class KinesisMessageDrivenChannelAdapterTests { .willReturn(new ListShardsResult() .withShards( new Shard().withShardId("closedShard1") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("1")))) + .withSequenceNumberRange(new SequenceNumberRange() + .withEndingSequenceNumber("10")))) .willReturn(new ListShardsResult() .withShards( new Shard().withShardId("closedShard1") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("1")), + .withSequenceNumberRange(new SequenceNumberRange() + .withEndingSequenceNumber("10")), new Shard().withShardId("newShard2") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("2")), + .withSequenceNumberRange(new SequenceNumberRange()), new Shard().withShardId("newShard3") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("3")), + .withSequenceNumberRange(new SequenceNumberRange()), new Shard().withShardId("closedShard4") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("4")))) + .withSequenceNumberRange(new SequenceNumberRange() + .withEndingSequenceNumber("40")), + new Shard().withShardId("closedEmptyShard5") + .withSequenceNumberRange(new SequenceNumberRange() + .withEndingSequenceNumber("50")))) .willReturn(new ListShardsResult() .withShards( new Shard().withShardId("closedShard1") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("1")), + .withSequenceNumberRange(new SequenceNumberRange() + .withEndingSequenceNumber("10")), new Shard().withShardId("newShard2") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("2")), + .withSequenceNumberRange(new SequenceNumberRange()), new Shard().withShardId("newShard3") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("3")), + .withSequenceNumberRange(new SequenceNumberRange()), new Shard().withShardId("closedShard4") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("4")), - new Shard().withShardId("newShard5") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("5")), + .withSequenceNumberRange(new SequenceNumberRange() + .withEndingSequenceNumber("40")), + new Shard().withShardId("closedEmptyShard5") + .withSequenceNumberRange(new SequenceNumberRange() + .withEndingSequenceNumber("50")), new Shard().withShardId("newShard6") - .withSequenceNumberRange(new SequenceNumberRange().withEndingSequenceNumber("6")))); + .withSequenceNumberRange(new SequenceNumberRange()), + new Shard().withShardId("newShard7") + .withSequenceNumberRange(new SequenceNumberRange()))); setClosedShard(amazonKinesis, "1"); setNewShard(amazonKinesis, "2"); setNewShard(amazonKinesis, "3"); setClosedShard(amazonKinesis, "4"); - setNewShard(amazonKinesis, "5"); + setClosedEmptyShard(amazonKinesis, "5"); setNewShard(amazonKinesis, "6"); + setNewShard(amazonKinesis, "7"); return amazonKinesis; } @@ -377,7 +397,8 @@ public class KinesisMessageDrivenChannelAdapterTests { String shardIterator = String.format("shard%sIterator1", shardIndex); given(amazonKinesis.getShardIterator( - KinesisShardOffset.latest(STREAM_FOR_RESHARDING, "closedShard" + shardIndex).toShardIteratorRequest())) + KinesisShardOffset.latest(STREAM_FOR_RESHARDING, "closedShard" + shardIndex) + .toShardIteratorRequest())) .willReturn(new GetShardIteratorResult().withShardIterator(shardIterator)); given(amazonKinesis.getRecords(new GetRecordsRequest().withShardIterator(shardIterator).withLimit(25))) @@ -386,6 +407,18 @@ public class KinesisMessageDrivenChannelAdapterTests { .withData(ByteBuffer.wrap("foo".getBytes())))); } + private void setClosedEmptyShard(AmazonKinesis amazonKinesis, String shardIndex) { + String shardIterator = String.format("shard%sIterator1", shardIndex); + + given(amazonKinesis.getShardIterator( + KinesisShardOffset.latest(STREAM_FOR_RESHARDING, "closedEmptyShard" + shardIndex) + .toShardIteratorRequest())) + .willReturn(new GetShardIteratorResult().withShardIterator(shardIterator)); + + given(amazonKinesis.getRecords(new GetRecordsRequest().withShardIterator(shardIterator).withLimit(25))) + .willReturn(new GetRecordsResult().withNextShardIterator(null)); + } + private void setNewShard(AmazonKinesis amazonKinesis, String shardIndex) { String shardIterator1 = String.format("shard%sIterator1", shardIndex); String shardIterator2 = String.format("shard%sIterator2", shardIndex); @@ -405,6 +438,11 @@ public class KinesisMessageDrivenChannelAdapterTests { .willReturn(new GetShardIteratorResult().withShardIterator(shardIterator2)); } + @Bean + public ConcurrentMetadataStore reshardingCheckpointStore() { + return new SimpleMetadataStore(); + } + @Bean public KinesisMessageDrivenChannelAdapter reshardingChannelAdapter() { KinesisMessageDrivenChannelAdapter adapter = new KinesisMessageDrivenChannelAdapter( @@ -415,6 +453,7 @@ public class KinesisMessageDrivenChannelAdapterTests { adapter.setDescribeStreamRetries(1); adapter.setRecordsLimit(25); adapter.setConcurrency(1); + adapter.setCheckpointStore(reshardingCheckpointStore()); DirectFieldAccessor dfa = new DirectFieldAccessor(adapter); dfa.setPropertyValue("describeStreamBackoff", 10);