From 4efb48d1a5aed55dfa7d2eaaa444ea95514f3be6 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 26 Sep 2024 12:06:18 -0400 Subject: [PATCH] GH-222: Expose KCL `gracefulShutdownTimeout` Fixes: #222 Issue link: https://github.com/spring-cloud/spring-cloud-stream-binder-aws-kinesis/issues/222 --- .../src/main/asciidoc/overview.adoc | 5 +++++ .../binder/kinesis/KinesisMessageChannelBinder.java | 1 + .../properties/KinesisConsumerProperties.java | 13 +++++++++++++ 3 files changed, 19 insertions(+) diff --git a/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc index a454696..73b1a5c 100644 --- a/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc +++ b/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc @@ -304,6 +304,11 @@ The KCL idle between requests in polling mode. + Default: 1500L +gracefulShutdownTimeout:: +The KCL graceful shutdown timeout in milliseconds. ++ +Default: 0 - regular shutdown process + Starting with version `4.0.4` (basically since `spring-integration-aws-3.0.8`), the `KclMessageDrivenChannelAdapter` can be customized programmatically for the `ConfigsBuilder` parts. For example, to set a custom value for the `LeaseManagementConfig.maxLeasesForWorker` property, the `ConsumerEndpointCustomizer` bean has to be provided: diff --git a/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/KinesisMessageChannelBinder.java b/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/KinesisMessageChannelBinder.java index 4785335..f9d7641 100644 --- a/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/KinesisMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/KinesisMessageChannelBinder.java @@ -417,6 +417,7 @@ public class KinesisMessageChannelBinder extends adapter.setLeaseTableName(kinesisConsumerProperties.getLeaseTableName()); adapter.setPollingMaxRecords(kinesisConsumerProperties.getPollingMaxRecords()); adapter.setPollingIdleTime(kinesisConsumerProperties.getPollingIdleTime()); + adapter.setGracefulShutdownTimeout(kinesisConsumerProperties.getGracefulShutdownTimeout()); if (properties.getExtension().isEmbedHeaders()) { adapter.setEmbeddedHeadersMapper(this.embeddedHeadersMapper); } diff --git a/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisConsumerProperties.java b/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisConsumerProperties.java index 0e08dda..2e8da02 100644 --- a/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisConsumerProperties.java +++ b/spring-cloud-stream-binder-kinesis/src/main/java/org/springframework/cloud/stream/binder/kinesis/properties/KinesisConsumerProperties.java @@ -88,6 +88,11 @@ public class KinesisConsumerProperties { */ private long pollingIdleTime = 1500L; + /** + * The KCL graceful shutdown timeout in milliseconds. + */ + private long gracefulShutdownTimeout; + private boolean embedHeaders; /** @@ -231,4 +236,12 @@ public class KinesisConsumerProperties { this.pollingIdleTime = pollingIdleTime; } + public long getGracefulShutdownTimeout() { + return this.gracefulShutdownTimeout; + } + + public void setGracefulShutdownTimeout(long gracefulShutdownTimeout) { + this.gracefulShutdownTimeout = gracefulShutdownTimeout; + } + }