[
https://issues.apache.org/jira/browse/BEAM-12459?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Boyuan Zhang updated BEAM-12459:
--------------------------------
Priority: P2 (was: P1)
> Watch does not properly advance the watermark by default
> --------------------------------------------------------
>
> Key: BEAM-12459
> URL: https://issues.apache.org/jira/browse/BEAM-12459
> Project: Beam
> Issue Type: Bug
> Components: sdk-java-core
> Reporter: Daniel Collins
> Assignee: Luke Cwik
> Priority: P2
>
> Assigning to Luke who has made substantial changes to this class most
> recently.
>
> It appears after investigation that when using Watch in the default
> configuration, the global watermark is not advanced properly, even though
> Watch documentation claims it should be
> ([https://github.com/apache/beam/blob/8922c1cf23c093262af9e4570d69947a9a749506/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Watch.java#L118]).
> The below is example code using Watch that will not advance the watermark,
> but should:
> ```
> Watch.growthOf(
> new PollFn<Integer, Integer>() {
> @Override
> public PollResult<Partition> apply(TopicPath element, Context c)
> { return PollResult.incomplete(Instant.now(), List.of(0)); }
> })
> .withPollInterval(...)
> .withTerminationPerInput(Watch.Growth.never());
> ```
> I've been advised that changing the return statement to `return
> PollResult.incomplete(Instant.now(),
> List.of(0)).withWatermark(Instant.now());` will resolve this issue, but the
> `withWatermark` function is commented as "By default, the watermark for a
> particular input is computed from a poll result as "earliest timestamp of new
> elements in this poll result". It can also be set explicitly via \{@link
> Growth.PollResult#withWatermark} if the \{@link Growth.PollFn} can provide a
> more optimistic estimate.". The goal is not to provide a more optimistic
> estimate, but to allow any advancement at all. If withWatermark is needed to
> close windows, this function should be required (or at least more prominent
> in all example code).
--
This message was sent by Atlassian Jira
(v8.3.4#803005)