riyarawat-amazon opened a new pull request, #258:
URL: https://github.com/apache/flink-connector-aws/pull/258
> **Note:** this depends on an AWS SDK v2 release whose DynamoDB Streams
model includes the
> `AT_TIMESTAMP` shard iterator type. Until that release CI fails to
compile; `aws.sdkv2.version`
> will be bumped in this PR once it is available.
### What is the purpose of the change
DynamoDB Streams now supports the `AT_TIMESTAMP` shard iterator type:
`GetShardIterator` with a
timestamp positions the iterator at the first record in the shard written at
or after that time.
This change adds `AT_TIMESTAMP` as an initial position for the DynamoDB
Streams source, so a job
can start reading from a point in time instead of only `LATEST` or
`TRIM_HORIZON`.
### Brief change log
- Add `InitialPosition.AT_TIMESTAMP` and the options
`flink.stream.initpos.timestamp` (date string
or epoch seconds) and `flink.stream.initpos.timestamp.format` (default
`yyyy-MM-dd'T'HH:mm:ss.SSSXXX`), matching the Kinesis source option names.
- Validate the configuration when the source is constructed: with
`AT_TIMESTAMP` the timestamp must
be set, parseable, and not in the future (a future timestamp is rejected
by `GetShardIterator`
and would otherwise make the job restart until the clock passes it).
- `SplitTracker`: under `AT_TIMESTAMP` every discovered shard — at
bootstrap, by periodic
discovery, and children discovered at shard end — starts at
`AT_TIMESTAMP(T)`. The service
resolves the position per shard, so a shard that straddles T never emits
records before T, and
a shard entirely before T finishes with no records. Ancestors of
already-tracked shards are
skipped; parent-before-child ordering is unchanged.
- `StartingPosition.atTimestamp(Instant)`; the proxy passes the timestamp on
`GetShardIteratorRequest`.
- Split serializer: version 3 stores the `AT_TIMESTAMP` instant; versions
0–2 remain readable.
- T is stored in the enumerator state, so on restore the configured
timestamp is not used.
- Docs: `AT_TIMESTAMP` in "Configuring Starting Position".
### Verifying this change
- Unit tests: config parsing/validation (missing, unparseable, epoch-seconds
fallback, future
timestamp), split tracker (positions at bootstrap / periodic discovery /
child handoff,
ancestor skipping, deduplication), enumerator (bootstrap lists all shards,
child gated behind
parent, child discovered only via handoff gets `AT_TIMESTAMP`), proxy
(timestamp passed on the
request) and serializer (round trip of the `AT_TIMESTAMP` instant, older
versions still read).
- Manually verified on a production stream (~33k shards) with a long-running
job cold-restarted
every 6 h at a random T within retention, plus scenario runs: T older than
retention behaves
like `TRIM_HORIZON`; changing T with a restored savepoint is ignored and
without state is
honoured; no record older than T was delivered and no per-key version gaps
were observed.
### Does this pull request potentially affect one of the following parts
- Dependencies (does it add or upgrade a dependency): **yes** — requires an
AWS SDK v2 version
whose DynamoDB Streams model includes `AT_TIMESTAMP` (see note above)
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: no (new config
options and enum value only)
- The serializers: **yes** — `DynamoDbStreamsShardSplitSerializer` version 3
(backward compatible)
- The runtime per-record code paths (performance sensitive): no
- Anything that affects deployment or recovery: JobManager (and its
components), Checkpointing,
Kubernetes/Yarn, ZooKeeper: no (enumerator state already carried the start
timestamp)
- The S3 file system connector: no
### Documentation
- Does this pull request introduce a new feature? yes
- If yes, how is the feature documented? docs (`dynamodb.md`)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]