Class DynamoDBStreamsDataFetcher

java.lang.Object
com.amazonaws.services.dynamodbv2.streamsadapter.DynamoDBStreamsDataFetcher
All Implemented Interfaces:
software.amazon.kinesis.retrieval.polling.DataFetcher

public class DynamoDBStreamsDataFetcher extends Object implements software.amazon.kinesis.retrieval.polling.DataFetcher
Implements fetching data from DynamoDB Streams using GetRecords and GetShardIterator API.
  • Field Summary

    Fields
    Modifier and Type
    Field
    Description
    protected static final int
     
  • Constructor Summary

    Constructors
    Constructor
    Description
    DynamoDBStreamsDataFetcher(@NotNull AmazonDynamoDBStreamsAdapterClient amazonDynamoDBStreamsAdapterClient, software.amazon.kinesis.retrieval.DataFetcherProviderConfig dynamoDBStreamsDataFetcherProviderConfig, @NotNull DynamoDBStreamsCatchUpConfig catchUpConfig)
     
  • Method Summary

    Modifier and Type
    Method
    Description
    void
    advanceIteratorTo(String sequenceNumber, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream)
    Advance the iterator to the given sequence number.
    software.amazon.kinesis.retrieval.GetRecordsResponseAdapter
    ddbGetRecords(@NonNull String nextIterator)
     
    software.amazon.awssdk.services.dynamodb.model.GetRecordsRequest
    Build the GetRecordsRequest with the given next iterator.
    protected software.amazon.awssdk.services.kinesis.model.DescribeStreamResponse
    getChildShards(String streamName, String shardId)
     
    software.amazon.awssdk.services.kinesis.model.GetRecordsRequest
     
    software.amazon.kinesis.retrieval.GetRecordsResponseAdapter
    getGetRecordsResponse(software.amazon.awssdk.services.dynamodb.model.GetRecordsRequest request)
    Call GetRecords API of DynamoDB Streams and return the result.
    software.amazon.awssdk.services.kinesis.model.GetRecordsResponse
    getGetRecordsResponse(software.amazon.awssdk.services.kinesis.model.GetRecordsRequest request)
     
    getNextIterator(software.amazon.awssdk.services.kinesis.model.GetShardIteratorRequest request)
    Get the next iterator using GetShardIterator API.
    software.amazon.kinesis.retrieval.DataFetcherResult
    Call GetRecords and get records back from DynamoDB Streams.
    software.amazon.awssdk.services.kinesis.model.GetRecordsResponse
    getRecords(@NonNull String nextIterator)
     
    void
    initialize(String initialCheckpoint, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream)
    Initialize the data fetcher with an initial checkpoint.
    void
    initialize(software.amazon.kinesis.retrieval.kpl.ExtendedSequenceNumber initialCheckpoint, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream)
     
    void
    resetIterator(String shardIterator, String sequenceNumber, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream)
    Reset the iterator to the given sharditerator, sequence number and initial position.
    void
    Restart the iterator using AT_SEQUENCE_NUMBER call.

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait

    Methods inherited from interface software.amazon.kinesis.retrieval.polling.DataFetcher

    getStreamIdentifier, isShardEndReached
  • Field Details

    • MAX_DESCRIBE_STREAM_ATTEMPTS_FOR_CHILD_SHARD_DISCOVERY_ON_NO_RECORDS

      protected static final int MAX_DESCRIBE_STREAM_ATTEMPTS_FOR_CHILD_SHARD_DISCOVERY_ON_NO_RECORDS
      See Also:
  • Constructor Details

    • DynamoDBStreamsDataFetcher

      public DynamoDBStreamsDataFetcher(@NotNull @NotNull AmazonDynamoDBStreamsAdapterClient amazonDynamoDBStreamsAdapterClient, software.amazon.kinesis.retrieval.DataFetcherProviderConfig dynamoDBStreamsDataFetcherProviderConfig, @NotNull @NotNull DynamoDBStreamsCatchUpConfig catchUpConfig)
  • Method Details

    • getRecords

      public software.amazon.kinesis.retrieval.DataFetcherResult getRecords()
      Call GetRecords and get records back from DynamoDB Streams.
      Specified by:
      getRecords in interface software.amazon.kinesis.retrieval.polling.DataFetcher
      Returns:
      DataFetcherResult containing records and whether the shard has reached end or not.
    • initialize

      public void initialize(String initialCheckpoint, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream)
      Initialize the data fetcher with an initial checkpoint.
      Specified by:
      initialize in interface software.amazon.kinesis.retrieval.polling.DataFetcher
      Parameters:
      initialCheckpoint - Current checkpoint sequence number for this shard. For TRIM_HORIZON and LATEST, it will be TRIM_HORIZON and LATEST. For everything else, it will be the last sequence number checkpointed in the lease table.
      initialPositionInStream - The initial position in stream. Will be either TRIM_HORIZON or LATEST for DynamoDB Streams.
    • initialize

      public void initialize(software.amazon.kinesis.retrieval.kpl.ExtendedSequenceNumber initialCheckpoint, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream)
      Specified by:
      initialize in interface software.amazon.kinesis.retrieval.polling.DataFetcher
    • advanceIteratorTo

      public void advanceIteratorTo(String sequenceNumber, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream)
      Advance the iterator to the given sequence number.
      Specified by:
      advanceIteratorTo in interface software.amazon.kinesis.retrieval.polling.DataFetcher
      Parameters:
      sequenceNumber - advance the iterator to the record at this sequence number.
      initialPositionInStream - The initialPositionInStream.
    • restartIterator

      public void restartIterator()
      Restart the iterator using AT_SEQUENCE_NUMBER call.
      Specified by:
      restartIterator in interface software.amazon.kinesis.retrieval.polling.DataFetcher
    • resetIterator

      public void resetIterator(String shardIterator, String sequenceNumber, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream)
      Reset the iterator to the given sharditerator, sequence number and initial position.
      Specified by:
      resetIterator in interface software.amazon.kinesis.retrieval.polling.DataFetcher
    • getGetRecordsResponse

      public software.amazon.awssdk.services.kinesis.model.GetRecordsResponse getGetRecordsResponse(software.amazon.awssdk.services.kinesis.model.GetRecordsRequest request) throws Exception
      Specified by:
      getGetRecordsResponse in interface software.amazon.kinesis.retrieval.polling.DataFetcher
      Throws:
      Exception
    • getGetRecordsRequest

      public software.amazon.awssdk.services.kinesis.model.GetRecordsRequest getGetRecordsRequest(String nextIterator)
      Specified by:
      getGetRecordsRequest in interface software.amazon.kinesis.retrieval.polling.DataFetcher
    • getGetRecordsResponse

      public software.amazon.kinesis.retrieval.GetRecordsResponseAdapter getGetRecordsResponse(software.amazon.awssdk.services.dynamodb.model.GetRecordsRequest request) throws ExecutionException, InterruptedException, TimeoutException
      Call GetRecords API of DynamoDB Streams and return the result.
      Parameters:
      request - the current get records request used to receive a response.
      Returns:
      GetRecordsResponse.
      Throws:
      ExecutionException
      InterruptedException
      TimeoutException
    • getChildShards

      protected software.amazon.awssdk.services.kinesis.model.DescribeStreamResponse getChildShards(String streamName, String shardId) throws InterruptedException
      Throws:
      InterruptedException
    • ddbGetRecordsRequest

      public software.amazon.awssdk.services.dynamodb.model.GetRecordsRequest ddbGetRecordsRequest(String nextIterator)
      Build the GetRecordsRequest with the given next iterator.
      Parameters:
      nextIterator - the next iterator to be used to build the GetRecordsRequest.
      Returns:
      GetRecordsRequest.
    • getNextIterator

      public String getNextIterator(software.amazon.awssdk.services.kinesis.model.GetShardIteratorRequest request) throws ExecutionException, InterruptedException, TimeoutException
      Get the next iterator using GetShardIterator API. If the shard is closed, return null. If the shard is not found, return null. If the shard is trimmed, call GetShardIterator API with TRIM_HORIZON.
      Specified by:
      getNextIterator in interface software.amazon.kinesis.retrieval.polling.DataFetcher
      Parameters:
      request - the current get shard iterator request used to receive a response.
      Returns:
      the next iterator.
      Throws:
      ExecutionException
      InterruptedException
      TimeoutException
    • getRecords

      public software.amazon.awssdk.services.kinesis.model.GetRecordsResponse getRecords(@NonNull @NonNull String nextIterator)
      Specified by:
      getRecords in interface software.amazon.kinesis.retrieval.polling.DataFetcher
    • ddbGetRecords

      public software.amazon.kinesis.retrieval.GetRecordsResponseAdapter ddbGetRecords(@NonNull @NonNull String nextIterator)