diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java index 63dd706..b0fc865 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java @@ -1535,7 +1535,12 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport lock.unlock(); } catch (Exception ex) { - logger.error(ex, () -> "Error during unlocking: " + lock); + if (KinesisMessageDrivenChannelAdapter.this.active) { + logger.error(ex, () -> "Error during unlocking: " + lock); + } + else { + logger.info(ex, () -> "Error during unlocking: " + lock + " while adapter was inactive"); + } } finally { iterator.remove(); diff --git a/src/main/java/org/springframework/integration/aws/outbound/KplMessageHandler.java b/src/main/java/org/springframework/integration/aws/outbound/KplMessageHandler.java index 67aa8d4..471741a 100644 --- a/src/main/java/org/springframework/integration/aws/outbound/KplMessageHandler.java +++ b/src/main/java/org/springframework/integration/aws/outbound/KplMessageHandler.java @@ -180,6 +180,16 @@ public class KplMessageHandler extends AbstractAwsMessageHandler implement this.embeddedHeadersMapper = embeddedHeadersMapper; } + /** + * Configure a {@link Duration} how often to call a {@link KinesisProducer#flush()}. + * @param flushDuration the {@link Duration} to periodic call of a {@link KinesisProducer#flush()}. + * @since 2.3.6 + */ + public void setFlushDuration(Duration flushDuration) { + Assert.notNull(flushDuration, "'flushDuration' must not be null."); + this.flushDuration = flushDuration; + } + /** * Unsupported operation. Use {@link #setEmbeddedHeadersMapper} instead. * @param headerMapper is not used.