From 59415106f30e04fc463331daa0e617d8d006823b Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 2 Oct 2023 17:07:50 -0400 Subject: [PATCH] GH-233: Add KCL polling config support Fixes https://github.com/spring-projects/spring-integration-aws/issues/233 --- .../KclMessageDrivenChannelAdapter.java | 46 +++++++++++++++---- 1 file changed, 38 insertions(+), 8 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 00b2782..18f44e8 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 @@ -40,7 +40,6 @@ 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; @@ -54,7 +53,9 @@ 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; import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisherAware; @@ -131,6 +132,8 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport private boolean bindSourceRecord; + private boolean fanOut = true; + private ApplicationEventPublisher applicationEventPublisher; private volatile Scheduler scheduler; @@ -248,6 +251,15 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport this.bindSourceRecord = bindSourceRecord; } + /** + * Specify a retrieval strategy: fan-out (true; default) or polling (false). + * @param fanOut false for a polling retrieval strategy. + * @since 3.0.2 + */ + public void setFanOut(boolean fanOut) { + this.fanOut = fanOut; + } + @Override protected void onInit() { super.onInit(); @@ -259,6 +271,28 @@ 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() { @@ -281,20 +315,16 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport + "because it does not make sense in case of [ListenerMode.batch]."); } - LifecycleConfig lifecycleConfig = this.config.lifecycleConfig().taskBackoffTimeMillis(this.consumerBackoff); - RetrievalConfig retrievalConfig = - this.config.retrievalConfig() - .glueSchemaRegistryDeserializer(this.glueSchemaRegistryDeserializer); this.scheduler = new Scheduler( this.config.checkpointConfig(), this.config.coordinatorConfig(), this.config.leaseManagementConfig(), - lifecycleConfig, + this.config.lifecycleConfig(), this.config.metricsConfig(), this.config.processorConfig(), - retrievalConfig); + this.config.retrievalConfig()); this.executor.execute(this.scheduler); }