[ 
https://issues.apache.org/jira/browse/FLINK-40500?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18123650#comment-18123650
 ] 

yongfu.gao commented on FLINK-40500:
------------------------------------

I'd like to pick this up.

I verified both claims on master: {{ExtractEventTimeProcessFunction}} builds 
the {{IdlenessTimer}} from the raw {{processingTimeService.getClock()}} 
(:90-92), and neither {{PausableRelativeClock}} nor a backpressure listener 
registration exists anywhere under {{{}flink-datastream{}}}, while the V1 paths 
do use them.

Reusing it should be clean: {{flink-datastream}} already imports 
{{org.apache.flink.runtime.*}} in production code, {{flink-runtime}} arrives 
transitively via {{{}flink-streaming-java{}}}, and there is no ArchUnit rule 
for {{{}flink-datastream{}}}. So I plan to follow V1 
({{{}TimestampsAndWatermarksOperator#open{}}}): wrap the clock in a 
{{{}PausableRelativeClock{}}}, use it for the {{{}IdlenessTimer{}}}, and 
register it as a backpressure listener.

 

> DataStream V2 idleness detection is not backpressure-aware - FLIP-471 was 
> never applied to ExtractEventTimeProcessFunction
> --------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40500
>                 URL: https://issues.apache.org/jira/browse/FLINK-40500
>             Project: Flink
>          Issue Type: Bug
>          Components: API / DataStream
>            Reporter: Martijn Visser
>            Priority: Major
>
> FLINK-35886 (FLIP-471) fixed incorrect idleness-timeout accounting when a 
> subtask is
> backpressured or blocked by watermark alignment, by introducing
> {{PausableRelativeClock}}. Every DataStream V1 path uses it
> ({{ProgressiveTimestampsAndWatermarks}}, {{TimestampsAndWatermarksOperator}},
> {{SourceOperator}}).
> The DataStream V2 {{ExtractEventTimeProcessFunction}} reuses
> {{WatermarksWithIdleness.IdlenessTimer}} but constructs it with the raw
> {{processingTimeService.getClock()}} 
> ({{ExtractEventTimeProcessFunction.java:90-92}});
> the hosting operators pass the service through unwrapped 
> ({{ProcessOperator.java:115-120}}
> and siblings). There is no {{PausableRelativeClock}} and no 
> backpressure-listener
> registration anywhere in flink-datastream (verified by grep). The V2 code 
> postdates the
> FLIP-471 fix by five months, so this is a missed carry-over, not a 
> merge-ordering issue.
> Consequence: under sustained backpressure, a V2 pipeline with an idle timeout 
> declares
> inputs idle while records are queued; the combined watermark advances past 
> them and they
> are dropped as late — exactly the failure FLIP-471 fixed for V1.
> Characterization test:
> {{ExtractEventTimeProcessFunctionTest#testIdleStatusEmittedPurelyOnWallClockElapse}}
>  —
> with idleTimeout=200ms, idle=true is emitted purely on wall-clock elapse; no 
> mechanism
> exists by which runtime-induced blocking could suppress it.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to