Class DynamoDBStreamsShardSyncer
java.lang.Object
software.amazon.kinesis.leases.HierarchicalShardSyncer
com.amazonaws.services.dynamodbv2.streamsadapter.DynamoDBStreamsShardSyncer
public class DynamoDBStreamsShardSyncer
extends software.amazon.kinesis.leases.HierarchicalShardSyncer
A shard syncer implementation specifically for DynamoDB Streams that extends the Kinesis hierarchical shard syncer.
This class handles synchronization of shard and lease states for DynamoDB streams in Lease Table
(Lease Management). It is responsible for ensuring that the leases table is up-to-date with the shards
in the stream.
-
Constructor Summary
ConstructorsConstructorDescriptionDynamoDBStreamsShardSyncer(boolean isMultiStreamMode, String streamIdentifier, boolean cleanupLeasesOfCompletedShards) Deprecated.DynamoDBStreamsShardSyncer(boolean isMultiStreamMode, String streamIdentifier, boolean cleanupLeasesOfCompletedShards, software.amazon.kinesis.coordinator.DeletedStreamListProvider deletedStreamListProvider) Deprecated.DynamoDBStreamsShardSyncer(boolean isMultiStreamMode, String streamIdentifier, boolean cleanupLeasesOfCompletedShards, software.amazon.kinesis.coordinator.DeletedStreamListProvider deletedStreamListProvider, software.amazon.kinesis.coordinator.StreamInfoManager streamInfoManager) Constructs a DynamoDBStreamsShardSyncer with support for deleted stream tracking.DynamoDBStreamsShardSyncer(boolean isMultiStreamMode, String streamIdentifier, boolean cleanupLeasesOfCompletedShards, software.amazon.kinesis.coordinator.StreamInfoManager streamInfoManager) Constructs a DynamoDBStreamsShardSyncer that can operate in either single or multi-stream mode. -
Method Summary
Modifier and TypeMethodDescriptionbooleancheckAndCreateLeaseForNewShards(@NonNull software.amazon.kinesis.leases.ShardDetector shardDetector, software.amazon.kinesis.leases.LeaseRefresher leaseRefresher, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPosition, software.amazon.kinesis.metrics.MetricsScope scope, boolean ignoreUnexpectedChildShards, boolean isLeaseTableEmpty) Checks for new shards and creates leases for them if they don't already exist.Methods inherited from class software.amazon.kinesis.leases.HierarchicalShardSyncer
checkAndCreateLeaseForNewShards, createLeaseForChildShard
-
Constructor Details
-
DynamoDBStreamsShardSyncer
@Deprecated public DynamoDBStreamsShardSyncer(boolean isMultiStreamMode, String streamIdentifier, boolean cleanupLeasesOfCompletedShards) Deprecated.Constructs a DynamoDBStreamsShardSyncer that can operate in either single or multi-stream mode.- Parameters:
isMultiStreamMode- Whether the syncer should operate in multi-stream modestreamIdentifier- The identifier for the stream being processedcleanupLeasesOfCompletedShards- Whether to cleanup the leases of finished shards.
-
DynamoDBStreamsShardSyncer
public DynamoDBStreamsShardSyncer(boolean isMultiStreamMode, String streamIdentifier, boolean cleanupLeasesOfCompletedShards, software.amazon.kinesis.coordinator.StreamInfoManager streamInfoManager) Constructs a DynamoDBStreamsShardSyncer that can operate in either single or multi-stream mode.- Parameters:
isMultiStreamMode- Whether the syncer should operate in multi-stream modestreamIdentifier- The identifier for the stream being processedcleanupLeasesOfCompletedShards- Whether to cleanup the leases of finished shards.streamInfoManager- Handles the lifecycle of stream metadata (e.g., streamId) in the CoordinatorState table.
-
DynamoDBStreamsShardSyncer
@Deprecated public DynamoDBStreamsShardSyncer(boolean isMultiStreamMode, String streamIdentifier, boolean cleanupLeasesOfCompletedShards, software.amazon.kinesis.coordinator.DeletedStreamListProvider deletedStreamListProvider) Deprecated.Constructs a DynamoDBStreamsShardSyncer with support for deleted stream tracking.- Parameters:
isMultiStreamMode- Whether the syncer should operate in multi-stream mode.streamIdentifier- The identifier for the stream being processed.deletedStreamListProvider- Provider for tracking deleted streams.cleanupLeasesOfCompletedShards- Whether to cleanup the leases of finished shards.
-
DynamoDBStreamsShardSyncer
public DynamoDBStreamsShardSyncer(boolean isMultiStreamMode, String streamIdentifier, boolean cleanupLeasesOfCompletedShards, software.amazon.kinesis.coordinator.DeletedStreamListProvider deletedStreamListProvider, software.amazon.kinesis.coordinator.StreamInfoManager streamInfoManager) Constructs a DynamoDBStreamsShardSyncer with support for deleted stream tracking.- Parameters:
isMultiStreamMode- Whether the syncer should operate in multi-stream mode.streamIdentifier- The identifier for the stream being processed.deletedStreamListProvider- Provider for tracking deleted streams.cleanupLeasesOfCompletedShards- Whether to cleanup the leases of finished shards.streamInfoManager- Handles the lifecycle of stream metadata (e.g., streamId) in the CoordinatorState table.
-
-
Method Details
-
checkAndCreateLeaseForNewShards
public boolean checkAndCreateLeaseForNewShards(@NonNull @NonNull software.amazon.kinesis.leases.ShardDetector shardDetector, software.amazon.kinesis.leases.LeaseRefresher leaseRefresher, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPosition, software.amazon.kinesis.metrics.MetricsScope scope, boolean ignoreUnexpectedChildShards, boolean isLeaseTableEmpty) throws software.amazon.kinesis.leases.exceptions.DependencyException, software.amazon.kinesis.leases.exceptions.InvalidStateException, software.amazon.kinesis.leases.exceptions.ProvisionedThroughputException, software.amazon.kinesis.exceptions.internal.KinesisClientLibIOException Checks for new shards and creates leases for them if they don't already exist. This method will: 1. Get the current stream shards 2. Sync the lease table with the current stream shards 3. Create leases for any new shards discovered- Overrides:
checkAndCreateLeaseForNewShardsin classsoftware.amazon.kinesis.leases.HierarchicalShardSyncer- Parameters:
shardDetector- Used to get information about the stream shardsleaseRefresher- Used to create and update leases in lease tableinitialPosition- The position where processing should start for newly discovered shardsscope- Metrics scope for recording operationsignoreUnexpectedChildShards- Whether to ignore child shards that violate assumptionsisLeaseTableEmpty- Whether the lease table is currently empty- Returns:
- true if the sync operation was successful
- Throws:
software.amazon.kinesis.leases.exceptions.DependencyException- If dependent services are unavailablesoftware.amazon.kinesis.leases.exceptions.InvalidStateException- If the lease table is in an invalid statesoftware.amazon.kinesis.leases.exceptions.ProvisionedThroughputException- If DynamoDB provisioned throughput is exceededsoftware.amazon.kinesis.exceptions.internal.KinesisClientLibIOException- If there are IO issues communicating with Kinesis
-