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 723d1de..675eb1c 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 @@ -289,6 +289,11 @@ Works only in `listenerMode.batch`. + Default: `false` +leaseTableName:: +The KCL table name for leases. ++ +Default: consumer group + 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 2db83f0..320bdc7 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 @@ -414,6 +414,7 @@ public class KinesisMessageChannelBinder extends adapter.setFanOut(kinesisConsumerProperties.isFanOut()); adapter.setMetricsLevel(kinesisConsumerProperties.getMetricsLevel()); adapter.setEmptyRecordList(kinesisConsumerProperties.isEmptyRecordList()); + adapter.setLeaseTableName(kinesisConsumerProperties.getLeaseTableName()); 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 c41269b..cef563a 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 @@ -72,6 +72,11 @@ public class KinesisConsumerProperties { */ private boolean emptyRecordList = false; + /** + * The KCL table name for leases. + */ + private String leaseTableName; + private boolean embedHeaders; /** @@ -191,4 +196,12 @@ public class KinesisConsumerProperties { this.emptyRecordList = emptyRecordList; } + public String getLeaseTableName() { + return this.leaseTableName; + } + + public void setLeaseTableName(String leaseTableName) { + this.leaseTableName = leaseTableName; + } + }