florianvazelle opened a new pull request, #29312:
URL: https://github.com/apache/flink/pull/29312
## What is the purpose of the change
A Python `TimestampAssigner` needs three operators: one extracts the
timestamp, one assigns the
timestamps and the watermarks, and one removes the extracted timestamp
again. Only the last one is
represented by the `DataStream` which `assign_timestamps_and_watermarks()`
returns, so a uid set on
that stream left the two others without one. With auto generated uids
disabled the job could then
not be submitted at all:
```
java.lang.IllegalStateException: Auto generated UIDs have been disabled but
no UID or hash has been
assigned to operator Timestamps/Watermarks
```
There was no way around it from user code, because the two operators are
only reachable through the
private `_j_data_stream`.
This pull request gives them a uid derived from the one the user sets, which
makes
`ExecutionConfig.disable_auto_generated_uids()` usable together with a
Python `TimestampAssigner`.
## Brief change log
- `assign_timestamps_and_watermarks()` returns a
`_TimestampsAndWatermarksDataStream`, which keeps
the transformations of the two operators it hides
- `uid()` on that stream also gives a uid to each of them, the one it
receives plus a fixed
suffix: `-extract-timestamp` and `-timestamps-and-watermarks`
- The suffixes are part of the uid, and so of the identity of the operator
in a savepoint. They
must stay as they are
- `name()` stays as it is: the names identify the operator in the metrics,
so a name set on the
stream must not change them
- The Java-only branch of `assign_timestamps_and_watermarks()`, which
builds one single operator,
is untouched
## Verifying this change
This change added tests and can be verified as follows:
- Added `CommonDataStreamTests.test_timestamp_assigner_uid`, which
disables auto generated uids,
builds a pipeline with a Python `TimestampAssigner` and a uid on every
operator, and asks for
the execution plan. Without the change it fails with the
`IllegalStateException` above
- The existing tests of `assign_timestamps_and_watermarks()` cover the
unchanged behaviour of the
three operators
## Does this pull request potentially affect one of the following parts:
- Dependencies (does it add or upgrade a dependency): **no**
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: **no** (no Java
class is changed; `DataStream.uid()` on the Python side does gain an
effect it did not have)
- The serializers: **no**
- The runtime per-record code paths (performance sensitive): **no**, the
change only builds the
stream graph
- Anything that affects deployment or recovery: JobManager (and its
components), Checkpointing,
Kubernetes/Yarn, ZooKeeper: **yes**. A job which uses a Python
`TimestampAssigner` *and* sets a
uid on the resulting stream gets a different operator ID for the two
hidden operators than
before, because the ID now comes from the uid instead of the topology.
Neither operator holds
user state, but a restore from an older savepoint may need
`--allowNonRestoredState`
- The S3 file system connector: **no**
## Documentation
- Does this pull request introduce a new feature? **no**, it is a bug fix
- If yes, how is the feature documented? **not applicable**
---
##### Was generative AI tooling used to co-author this PR?
- [X] Yes (please specify the tool below)
Generated-by: Claude Code (Claude Opus 5)
--
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]