Florian Vazelle created FLINK-40829:
---------------------------------------
Summary: operators hidden by `assign_timestamps_and_watermarks`
get no uid, breaking `pipeline.auto-generate-uids: false`
Key: FLINK-40829
URL: https://issues.apache.org/jira/browse/FLINK-40829
Project: Flink
Issue Type: Bug
Components: API / Python
Reporter: Florian Vazelle
{{DataStream.assign_timestamps_and_watermarks()}} expands a *Python*
{{TimestampAssigner}} into three transformations, but returns a {{DataStream}}
which represents only the last one:
{noformat}
Extract-Timestamp (python) -> Timestamps/Watermarks (java) ->
Remove-Timestamp (java map){noformat}
A uid set on the returned stream reaches {{Remove-Timestamp}} only. The two
others are not reachable from the Python API, and with
{{pipeline.auto-generate-uids: false}} (or
{{{}ExecutionConfig.disable_auto_generated_uids(){}}}) the job cannot be
submitted at all:
{noformat}
java.lang.IllegalStateException: Auto generated UIDs have been disabled but no
UID or hash has been assigned to operator Timestamps/Watermarks
{noformat}
h2. How to reproduce
{code:python}
from pyflink.common import WatermarkStrategy
from pyflink.common.watermark_strategy import TimestampAssigner
from pyflink.datastream import StreamExecutionEnvironment
class MyAssigner(TimestampAssigner):
def extract_timestamp(self, value, record_timestamp):
return value[0]
env = StreamExecutionEnvironment.get_execution_environment()
env.get_config().disable_auto_generated_uids()
stream = env.from_collection([(1000, "a"), (2000, "b")]).uid("source")
timed = stream.assign_timestamps_and_watermarks(
WatermarkStrategy.for_monotonous_timestamps().with_timestamp_assigner(MyAssigner())
).uid("watermarks")
timed.print().uid("sink")
env.execute()
{code}
*Expected:* the job is submitted, because the uid the user sets covers every
operator the call created.
*Actual:* submission fails on {{{}Extract-Timestamp{}}}, and on
{{Timestamps/Watermarks}} once that one has a uid.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)