[ https://issues.apache.org/jira/browse/FLINK-31373?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel ]
Flink Jira Bot updated FLINK-31373: ----------------------------------- Labels: pull-request-available stale-major (was: pull-request-available) I am the [Flink Jira Bot|https://github.com/apache/flink-jira-bot/] and I help the community manage its development. I see this issues has been marked as Major but is unassigned and neither itself nor its Sub-Tasks have been updated for 60 days. I have gone ahead and added a "stale-major" to the issue". If this ticket is a Major, please either assign yourself or give an update. Afterwards, please remove the label or in 7 days the issue will be deprioritized. > PerRoundWrapperOperator should carry epoch information in watermark > ------------------------------------------------------------------- > > Key: FLINK-31373 > URL: https://issues.apache.org/jira/browse/FLINK-31373 > Project: Flink > Issue Type: Bug > Components: Library / Machine Learning > Affects Versions: ml-2.2.0 > Reporter: Zhipeng Zhang > Priority: Major > Labels: pull-request-available, stale-major > > Currently we use PerRoundWrapperOperator to wrap the normal flink operators > such that they can be used in iterations. > We already contained the epoch information in each record so that we know > which iteration each record belongs to. > However, there is no epoch information when the stream element is a > watermark. This works in most cases, but fail to address the following use > case: > - In DataStreamUtils#withBroadcast, we will cache the elements (including > watermarks) from non-broadcast inputs until the broadcast variables are > ready. When the broadcast variables are ready, once we receive a stream > element we will process the cached elements first. If the received element is > a watermark, the current implementation of iteration module fails > (ProxyOutput#collect throws NPE) since there is no epoch information. -- This message was sent by Atlassian Jira (v8.20.10#820010)