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
FieldsModifier and TypeFieldDescriptionprotected static final int -
Constructor Summary
ConstructorsConstructorDescriptionDynamoDBStreamsDataFetcher(@NotNull AmazonDynamoDBStreamsAdapterClient amazonDynamoDBStreamsAdapterClient, software.amazon.kinesis.retrieval.DataFetcherProviderConfig dynamoDBStreamsDataFetcherProviderConfig, @NotNull DynamoDBStreamsCatchUpConfig catchUpConfig) -
Method Summary
Modifier and TypeMethodDescriptionvoidadvanceIteratorTo(String sequenceNumber, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream) Advance the iterator to the given sequence number.software.amazon.kinesis.retrieval.GetRecordsResponseAdapterddbGetRecords(@NonNull String nextIterator) software.amazon.awssdk.services.dynamodb.model.GetRecordsRequestddbGetRecordsRequest(String nextIterator) Build the GetRecordsRequest with the given next iterator.protected software.amazon.awssdk.services.kinesis.model.DescribeStreamResponsegetChildShards(String streamName, String shardId) software.amazon.awssdk.services.kinesis.model.GetRecordsRequestgetGetRecordsRequest(String nextIterator) software.amazon.kinesis.retrieval.GetRecordsResponseAdaptergetGetRecordsResponse(software.amazon.awssdk.services.dynamodb.model.GetRecordsRequest request) Call GetRecords API of DynamoDB Streams and return the result.software.amazon.awssdk.services.kinesis.model.GetRecordsResponsegetGetRecordsResponse(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.DataFetcherResultCall GetRecords and get records back from DynamoDB Streams.software.amazon.awssdk.services.kinesis.model.GetRecordsResponsegetRecords(@NonNull String nextIterator) voidinitialize(String initialCheckpoint, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream) Initialize the data fetcher with an initial checkpoint.voidinitialize(software.amazon.kinesis.retrieval.kpl.ExtendedSequenceNumber initialCheckpoint, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream) voidresetIterator(String shardIterator, String sequenceNumber, software.amazon.kinesis.common.InitialPositionInStreamExtended initialPositionInStream) Reset the iterator to the given sharditerator, sequence number and initial position.voidRestart the iterator using AT_SEQUENCE_NUMBER call.Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods 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:
getRecordsin interfacesoftware.amazon.kinesis.retrieval.polling.DataFetcher- Returns:
DataFetcherResultcontaining 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:
initializein interfacesoftware.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:
initializein interfacesoftware.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:
advanceIteratorToin interfacesoftware.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:
restartIteratorin interfacesoftware.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:
resetIteratorin interfacesoftware.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:
getGetRecordsResponsein interfacesoftware.amazon.kinesis.retrieval.polling.DataFetcher- Throws:
Exception
-
getGetRecordsRequest
public software.amazon.awssdk.services.kinesis.model.GetRecordsRequest getGetRecordsRequest(String nextIterator) - Specified by:
getGetRecordsRequestin interfacesoftware.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:
ExecutionExceptionInterruptedExceptionTimeoutException
-
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:
getNextIteratorin interfacesoftware.amazon.kinesis.retrieval.polling.DataFetcher- Parameters:
request- the current get shard iterator request used to receive a response.- Returns:
- the next iterator.
- Throws:
ExecutionExceptionInterruptedExceptionTimeoutException
-
getRecords
public software.amazon.awssdk.services.kinesis.model.GetRecordsResponse getRecords(@NonNull @NonNull String nextIterator) - Specified by:
getRecordsin interfacesoftware.amazon.kinesis.retrieval.polling.DataFetcher
-
ddbGetRecords
public software.amazon.kinesis.retrieval.GetRecordsResponseAdapter ddbGetRecords(@NonNull @NonNull String nextIterator)
-