rkhachatryan opened a new pull request, #29270:
URL: https://github.com/apache/flink/pull/29270

   A source subtask can announce WatermarkStatus.IDLE downstream and never take 
it back while it keeps emitting records. Downstream StatusWatermarkValve then 
excludes the channel, stalling the watermark at parallelism 1 and dropping the 
records as late at higher parallelism.
   
   Three places on the path from a split's watermark generator to the data 
output treat an advancing watermark as the only evidence of activity:
   
     - WatermarkOutputMultiplexer.updateCombinedWatermark() forwards markIdle() 
when the combined status is idle but nothing when it is active, so an output 
that is active without advancing its watermark produces neither an 
emitWatermark nor a markIdle call. Because the split branch of 
ProgressiveTimestampsAndWatermarks.IdlenessManager starts out idle and is only 
ever cleared by one of those two calls, an active split is then ANDed into IDLE 
by the main output's idleness timer .
     - WatermarksWithIdleness.onEvent() resets the idleness timer but does not 
undo the markIdle() it emitted earlier .
     - WatermarkToDataOutput.emitWatermark() gates markActiveInternally() on 
the monotonicity check, so a non-advancing watermark is swallowed together with 
the activation .
   
   They compound: once all splits fall idle the combined watermark is flushed 
to the maximum over all outputs, so a split that resumes behind that maximum 
produces no advancing watermark by construction.
   
   Propagate the activity in all three, guarding the multiplexer on the new 
CombinedWatermarkStatus.hasOutputs(). That guard is load-bearing: 
updateCombinedWatermark() returns early on an empty output list and leaves the 
idle flag uncomputed, and a subtask that was never assigned a split has to be 
able to stand down via its main-output idleness timer alone.
   
   <!--
   *Thank you very much for contributing to Apache Flink - we are happy that 
you want to help us improve Flink. To help the community review your 
contribution in the best possible way, please go through the checklist below, 
which will get the contribution into a shape in which it can be best reviewed.*
   
   *Please understand that we do not do this to make contributions to Flink a 
hassle. In order to uphold a high standard of quality for code contributions, 
while at the same time managing a large number of contributions, we need 
contributors to prepare the contributions well, and give reviewers enough 
contextual information for the review. Please also understand that 
contributions that do not follow this guide will take longer to review and thus 
typically be picked up with lower priority by the community.*
   
   ## Contribution Checklist
   
     - Make sure that the pull request corresponds to a [JIRA 
issue](https://issues.apache.org/jira/projects/FLINK/issues). Exceptions are 
made for typos in JavaDoc or documentation files, which need no JIRA issue.
     
     - Name the pull request in the form "[FLINK-XXXX] [component] Title of the 
pull request", where *FLINK-XXXX* should be replaced by the actual issue 
number. Skip *component* if you are unsure about which is the best component.
     Typo fixes that have no associated JIRA issue should be named following 
this pattern: `[hotfix] [docs] Fix typo in event time introduction` or 
`[hotfix] [javadocs] Expand JavaDoc for PuncuatedWatermarkGenerator`.
   
     - Fill out the template below to describe the changes contributed by the 
pull request. That will give reviewers the context they need to do the review.
     
     - Make sure that the change passes the automated tests, i.e., `mvn clean 
verify` passes. You can set up Azure Pipelines CI to do that following [this 
guide](https://cwiki.apache.org/confluence/display/FLINK/Azure+Pipelines#AzurePipelines-Tutorial:SettingupAzurePipelinesforaforkoftheFlinkrepository).
   
     - Each pull request should address only one issue, not mix up code from 
multiple issues.
     
     - Each commit in the pull request has a meaningful commit message 
(including the JIRA id)
   
     - Once all items of the checklist are addressed, remove the above text and 
this checklist, leaving only the filled out template below.
   
   
   **(The sections below can be removed for hotfixes of typos)**
   -->
   
   


-- 
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]

Reply via email to