[ 
https://issues.apache.org/jira/browse/FLINK-40829?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated FLINK-40829:
-----------------------------------
    Labels: pull-request-available  (was: )

> 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
>            Priority: Minor
>              Labels: pull-request-available
>
> {{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)

Reply via email to