From 89b19f1e7ac5d0fa900e87e18dae3b4b764c7aac Mon Sep 17 00:00:00 2001 From: siddharthjain210 Date: Mon, 30 Oct 2023 20:34:24 +0530 Subject: [PATCH] GH-235: Fix Retrieval & Lifecycle Config for KCL Fixes: https://github.com/spring-projects/spring-integration-aws/issues/235 * Corrected the initialization of Retrieval and Lifecycle Config. The start Scheduler was creating different instances. * GH-235 Fixing the failing tests for Enhanced Fan Out. --- .../KclMessageDrivenChannelAdapter.java | 49 ++++++++++--------- .../KclMessageDrivenChannelAdapterTests.java | 3 +- 2 files changed, 27 insertions(+), 25 deletions(-) diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java index 224ce20..a1479cd 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java @@ -44,6 +44,7 @@ import software.amazon.kinesis.coordinator.Scheduler; import software.amazon.kinesis.exceptions.InvalidStateException; import software.amazon.kinesis.exceptions.ShutdownException; import software.amazon.kinesis.exceptions.ThrottlingException; +import software.amazon.kinesis.lifecycle.LifecycleConfig; import software.amazon.kinesis.lifecycle.events.InitializationInput; import software.amazon.kinesis.lifecycle.events.LeaseLostInput; import software.amazon.kinesis.lifecycle.events.ProcessRecordsInput; @@ -57,6 +58,7 @@ import software.amazon.kinesis.processor.ShardRecordProcessorFactory; import software.amazon.kinesis.processor.SingleStreamTracker; import software.amazon.kinesis.processor.StreamTracker; import software.amazon.kinesis.retrieval.KinesisClientRecord; +import software.amazon.kinesis.retrieval.RetrievalConfig; import software.amazon.kinesis.retrieval.RetrievalSpecificConfig; import software.amazon.kinesis.retrieval.fanout.FanOutConfig; import software.amazon.kinesis.retrieval.polling.PollingConfig; @@ -276,28 +278,6 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport this.cloudWatchClient, this.workerId, this.recordProcessorFactory); - - this.config.lifecycleConfig().taskBackoffTimeMillis(this.consumerBackoff); - - RetrievalSpecificConfig retrievalSpecificConfig; - - String singleStreamName = this.streams.length == 1 ? this.streams[0] : null; - - if (this.fanOut) { - retrievalSpecificConfig = - new FanOutConfig(this.kinesisClient) - .applicationName(this.consumerGroup) - .streamName(singleStreamName); - } - else { - retrievalSpecificConfig = - new PollingConfig(this.kinesisClient) - .streamName(singleStreamName); - } - - this.config.retrievalConfig() - .glueSchemaRegistryDeserializer(this.glueSchemaRegistryDeserializer) - .retrievalSpecificConfig(retrievalSpecificConfig); } private StreamTracker buildStreamTracker() { @@ -320,15 +300,36 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport + "because it does not make sense in case of [ListenerMode.batch]."); } + LifecycleConfig lifecycleConfig = this.config.lifecycleConfig(); + lifecycleConfig.taskBackoffTimeMillis(this.consumerBackoff); + + RetrievalSpecificConfig retrievalSpecificConfig; + String singleStreamName = this.streams.length == 1 ? this.streams[0] : null; + if (this.fanOut) { + retrievalSpecificConfig = + new FanOutConfig(this.kinesisClient) + .applicationName(this.consumerGroup) + .streamName(singleStreamName); + } + else { + retrievalSpecificConfig = + new PollingConfig(this.kinesisClient) + .streamName(singleStreamName); + } + + RetrievalConfig retrievalConfig = this.config.retrievalConfig() + .glueSchemaRegistryDeserializer(this.glueSchemaRegistryDeserializer) + .retrievalSpecificConfig(retrievalSpecificConfig); + this.scheduler = new Scheduler( this.config.checkpointConfig(), this.config.coordinatorConfig(), this.config.leaseManagementConfig(), - this.config.lifecycleConfig(), + lifecycleConfig, this.config.metricsConfig(), this.config.processorConfig(), - this.config.retrievalConfig()); + retrievalConfig); this.executor.execute(this.scheduler); } diff --git a/src/test/java/org/springframework/integration/aws/kinesis/KclMessageDrivenChannelAdapterTests.java b/src/test/java/org/springframework/integration/aws/kinesis/KclMessageDrivenChannelAdapterTests.java index 80c1b00..65f498c 100644 --- a/src/test/java/org/springframework/integration/aws/kinesis/KclMessageDrivenChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/kinesis/KclMessageDrivenChannelAdapterTests.java @@ -112,7 +112,8 @@ public class KclMessageDrivenChannelAdapterTests implements LocalstackContainerT .join() .consumers(); - assertThat(streamConsumers).hasSize(1); + // Because FanOut is false, there would be no Stream Consumers. + assertThat(streamConsumers).hasSize(0); } @Configuration