Add lockRenewalTimeout to the KinesisMesChAd

Related to https://github.com/spring-cloud/spring-cloud-stream-binder-aws-kinesis/issues/148
This commit is contained in:
Artem Bilan
2021-01-13 15:30:07 -05:00
parent b43a230ffe
commit ea135917f3

View File

@@ -161,6 +161,8 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport
private int describeStreamRetries = 50;
private long lockRenewalTimeout = 10_000L;
private boolean resetCheckpoints;
private InboundMessageMapper<byte[]> embeddedHeadersMapper;
@@ -291,6 +293,16 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport
this.startTimeout = startTimeout;
}
/**
* Configure a timeout in milliseconds to wait for lock on shard renewal.
* @param lockRenewalTimeout the timeout to wait for lock renew in milliseconds.
* @since 2.3.5
*/
public void setLockRenewalTimeout(long lockRenewalTimeout) {
Assert.isTrue(lockRenewalTimeout > 0, "'lockRenewalTimeout' must be more than 0");
this.lockRenewalTimeout = lockRenewalTimeout;
}
/**
* The maximum number of concurrent {@link ConsumerInvoker}s running. The {@link ShardConsumer}s
* are evenly distributed between {@link ConsumerInvoker}s. Messages from within the same shard
@@ -920,7 +932,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport
LockCompletableFuture unlockFuture = new LockCompletableFuture(this.key);
KinesisMessageDrivenChannelAdapter.this.shardConsumerManager.unlock(unlockFuture);
try {
unlockFuture.get(1, TimeUnit.SECONDS);
unlockFuture.get(KinesisMessageDrivenChannelAdapter.this.lockRenewalTimeout, TimeUnit.MILLISECONDS);
}
catch (Exception ex) {
if (ex instanceof InterruptedException) {
@@ -1029,7 +1041,8 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport
KinesisMessageDrivenChannelAdapter.this.shardConsumerManager.renewLock(renewLockFuture);
boolean lockRenewed = false;
try {
lockRenewed = renewLockFuture.get(1, TimeUnit.SECONDS);
lockRenewed = renewLockFuture.get(KinesisMessageDrivenChannelAdapter.this.lockRenewalTimeout,
TimeUnit.MILLISECONDS);
}
catch (Exception ex) {
if (ex instanceof InterruptedException) {