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 530cdbd..dccab8c 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 @@ -1557,8 +1557,13 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport try { lock.unlock(); } - catch (Exception e) { - logger.error("Error during unlocking: " + lock, e); + catch (Exception ex) { + if (KinesisMessageDrivenChannelAdapter.this.active) { + logger.error("Error during unlocking: " + lock, ex); + } + else { + logger.info("Error during unlocking: " + lock + " while adapter was inactive", ex); + } } 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.