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)

Reply via email to