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]