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

    Constructors
    Constructor
    Description
    DynamoDBStreamsShardSyncer(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 Type
    Method
    Description
    boolean
    checkAndCreateLeaseForNewShards(@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

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • 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 mode
      streamIdentifier - The identifier for the stream being processed
      cleanupLeasesOfCompletedShards - 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 mode
      streamIdentifier - The identifier for the stream being processed
      cleanupLeasesOfCompletedShards - 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:
      checkAndCreateLeaseForNewShards in class software.amazon.kinesis.leases.HierarchicalShardSyncer
      Parameters:
      shardDetector - Used to get information about the stream shards
      leaseRefresher - Used to create and update leases in lease table
      initialPosition - The position where processing should start for newly discovered shards
      scope - Metrics scope for recording operations
      ignoreUnexpectedChildShards - Whether to ignore child shards that violate assumptions
      isLeaseTableEmpty - 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 unavailable
      software.amazon.kinesis.leases.exceptions.InvalidStateException - If the lease table is in an invalid state
      software.amazon.kinesis.leases.exceptions.ProvisionedThroughputException - If DynamoDB provisioned throughput is exceeded
      software.amazon.kinesis.exceptions.internal.KinesisClientLibIOException - If there are IO issues communicating with Kinesis