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
This commit is contained in:
@@ -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<Shard>, List<Shard>> 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(
|
||||
|
||||
@@ -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<KinesisShardOffset> 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);
|
||||
|
||||
Reference in New Issue
Block a user