Allow specifying a custom LeaseManager in Worker.Builder
This commit is contained in:
parent
147302b398
commit
1f7a05eb02
1 changed files with 19 additions and 2 deletions
|
|
@ -33,6 +33,8 @@ import java.util.concurrent.TimeUnit;
|
||||||
import java.util.concurrent.TimeoutException;
|
import java.util.concurrent.TimeoutException;
|
||||||
|
|
||||||
import com.amazonaws.services.kinesis.clientlibrary.proxies.IKinesisProxy;
|
import com.amazonaws.services.kinesis.clientlibrary.proxies.IKinesisProxy;
|
||||||
|
import com.amazonaws.services.kinesis.leases.interfaces.ILeaseManager;
|
||||||
|
|
||||||
import org.apache.commons.logging.Log;
|
import org.apache.commons.logging.Log;
|
||||||
import org.apache.commons.logging.LogFactory;
|
import org.apache.commons.logging.LogFactory;
|
||||||
|
|
||||||
|
|
@ -1076,6 +1078,7 @@ public class Worker implements Runnable {
|
||||||
private AmazonDynamoDB dynamoDBClient;
|
private AmazonDynamoDB dynamoDBClient;
|
||||||
private AmazonCloudWatch cloudWatchClient;
|
private AmazonCloudWatch cloudWatchClient;
|
||||||
private IMetricsFactory metricsFactory;
|
private IMetricsFactory metricsFactory;
|
||||||
|
private ILeaseManager<KinesisClientLease> leaseManager;
|
||||||
private ExecutorService execService;
|
private ExecutorService execService;
|
||||||
private ShardPrioritization shardPrioritization;
|
private ShardPrioritization shardPrioritization;
|
||||||
private IKinesisProxy kinesisProxy;
|
private IKinesisProxy kinesisProxy;
|
||||||
|
|
@ -1173,6 +1176,18 @@ public class Worker implements Runnable {
|
||||||
return this;
|
return this;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Set the lease manager.
|
||||||
|
*
|
||||||
|
* @param leaseManager
|
||||||
|
* Lease manager used to manage shard leases
|
||||||
|
* @return A reference to this updated object so that method calls can be chained together.
|
||||||
|
*/
|
||||||
|
public Builder leaseManager(ILeaseManager<KinesisClientLease> leaseManager) {
|
||||||
|
this.leaseManager = leaseManager;
|
||||||
|
return this;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Set the executor service for processing records.
|
* Set the executor service for processing records.
|
||||||
*
|
*
|
||||||
|
|
@ -1273,6 +1288,9 @@ public class Worker implements Runnable {
|
||||||
if (metricsFactory == null) {
|
if (metricsFactory == null) {
|
||||||
metricsFactory = getMetricsFactory(cloudWatchClient, config);
|
metricsFactory = getMetricsFactory(cloudWatchClient, config);
|
||||||
}
|
}
|
||||||
|
if (leaseManager == null) {
|
||||||
|
leaseManager = new KinesisClientLeaseManager(config.getTableName(), dynamoDBClient);
|
||||||
|
}
|
||||||
if (shardPrioritization == null) {
|
if (shardPrioritization == null) {
|
||||||
shardPrioritization = new ParentsFirstShardPrioritization(1);
|
shardPrioritization = new ParentsFirstShardPrioritization(1);
|
||||||
}
|
}
|
||||||
|
|
@ -1294,8 +1312,7 @@ public class Worker implements Runnable {
|
||||||
config.getShardSyncIntervalMillis(),
|
config.getShardSyncIntervalMillis(),
|
||||||
config.shouldCleanupLeasesUponShardCompletion(),
|
config.shouldCleanupLeasesUponShardCompletion(),
|
||||||
null,
|
null,
|
||||||
new KinesisClientLibLeaseCoordinator(new KinesisClientLeaseManager(config.getTableName(),
|
new KinesisClientLibLeaseCoordinator(leaseManager,
|
||||||
dynamoDBClient),
|
|
||||||
config.getWorkerIdentifier(),
|
config.getWorkerIdentifier(),
|
||||||
config.getFailoverTimeMillis(),
|
config.getFailoverTimeMillis(),
|
||||||
config.getEpsilonMillis(),
|
config.getEpsilonMillis(),
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue