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 35f4986..3026ef9 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 @@ -1,5 +1,5 @@ /* - * Copyright 2017-2022 the original author or authors. + * Copyright 2017-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -1561,7 +1561,7 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport } } - sleep(250, + sleep(1000, new IllegalStateException("ShardConsumerManager Thread [" + this + "] has been interrupted"), true); } diff --git a/src/main/java/org/springframework/integration/aws/lock/DynamoDbLockRegistry.java b/src/main/java/org/springframework/integration/aws/lock/DynamoDbLockRegistry.java index fcbdda6..5b8fe5e 100644 --- a/src/main/java/org/springframework/integration/aws/lock/DynamoDbLockRegistry.java +++ b/src/main/java/org/springframework/integration/aws/lock/DynamoDbLockRegistry.java @@ -26,7 +26,6 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executor; -import java.util.concurrent.ExecutorService; import java.util.concurrent.ThreadFactory; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Condition; @@ -137,12 +136,6 @@ public class DynamoDbLockRegistry implements ExpirableLockRegistry, Initializing private long heartbeatPeriod = 5L; - /** - * Flag to denote whether the {@link ExecutorService} was provided via the setter and - * thus should not be shutdown when {@link #destroy()} is called. - */ - private boolean executorExplicitlySet; - private volatile boolean initialized; public DynamoDbLockRegistry(AmazonDynamoDB dynamoDB) { @@ -204,6 +197,11 @@ public class DynamoDbLockRegistry implements ExpirableLockRegistry, Initializing this.leaseDuration = leaseDuration; } + /** + * Specify a period in milliseconds how often send locks renewal requests called heartbeat. + * When the value is less than or equal to {@code 0}, the heartbeat is disabled. + * @param heartbeatPeriod the heartbeat period for background thread to renew locks in DB + */ public void setHeartbeatPeriod(long heartbeatPeriod) { this.heartbeatPeriod = heartbeatPeriod; } @@ -227,7 +225,9 @@ public class DynamoDbLockRegistry implements ExpirableLockRegistry, Initializing if (!this.dynamoDBLockClientExplicitlySet) { AmazonDynamoDBLockClientOptions dynamoDBLockClientOptions = AmazonDynamoDBLockClientOptions .builder(this.dynamoDB, this.tableName).withPartitionKeyName(this.partitionKey) - .withSortKeyName(this.sortKeyName).withHeartbeatPeriod(this.heartbeatPeriod) + .withSortKeyName(this.sortKeyName) + .withCreateHeartbeatBackgroundThread(this.heartbeatPeriod > 0) + .withHeartbeatPeriod(this.heartbeatPeriod) .withLeaseDuration(this.leaseDuration).build(); this.dynamoDBLockClient = new AmazonDynamoDBLockClient(dynamoDBLockClientOptions);