From 526353c9d80c6234e42fbf06ae84f81b77a4f7eb Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 21 Dec 2020 13:24:52 -0500 Subject: [PATCH] Kinesis CA: Fulfill `lockFuture` on exception Related to https://github.com/spring-cloud/spring-cloud-stream-binder-aws-kinesis/issues/148 When `lock.tryLock()` ends up with an exception, we just log it under error category. * Add also `lockFuture.complete(false)` in the catch block when we try to renew the lock --- .../aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java | 1 + 1 file changed, 1 insertion(+) 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 c398595..e5fbb66 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 @@ -1506,6 +1506,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport } } catch (Exception e) { + lockFuture.complete(false); logger.error("Error during locking: " + lock, e); } }