sjvanrossum commented on code in PR #40176:
URL: https://github.com/apache/beam/pull/40176#discussion_r4061460082


##########
sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/CustomTimestampPolicyWithLimitedDelay.java:
##########
@@ -86,11 +91,27 @@ public Instant getWatermark(PartitionContext ctx) {
 
   @VisibleForTesting
   Instant getWatermark(PartitionContext ctx, Instant now) {
+    // The watermark must not move backwards. The idle branch answers from 
'backlogCheckTime' while
+    // the fallback answers from 'maxEventTimestamp', and only the latter is 
advanced by records, so
+    // without this clamp a partition which advanced while idle would regress 
all the way back to
+    // 'maxEventTimestamp - maxDelay' the moment its first record arrived.
+    Instant candidate = candidateWatermark(ctx, now);
+    if (candidate.isAfter(lastWatermark)) {
+      lastWatermark = candidate;
+    }
+    return lastWatermark;
+  }
+
+  private Instant candidateWatermark(PartitionContext ctx, Instant now) {
     if (maxEventTimestamp.isAfter(now)) {
       return now.minus(maxDelay); // (a) above.
     } else if (ctx.getMessageBacklog() == 0
-        && 
ctx.getBacklogCheckTime().minus(maxDelay).isAfter(maxEventTimestamp) // Idle
-        && maxEventTimestamp.getMillis() > 0) { // Read at least one record 
with positive timestamp.
+        && 
ctx.getBacklogCheckTime().minus(maxDelay).isAfter(maxEventTimestamp)) { // Idle
+      // A zero backlog means the reader has a position and knows it is at the 
log end, so no
+      // unread record can arrive late regardless of whether one has ever been 
read. Requiring a
+      // record to have been read here as well would pin a partition which is 
caught up but has
+      // delivered nothing since the job started at 'maxEventTimestamp - 
maxDelay' forever, because
+      // only a delivered record can advance 'maxEventTimestamp'.
       return ctx.getBacklogCheckTime().minus(maxDelay);

Review Comment:
   After looking at this some more it seems like the intent was to advance the 
watermark unconditionally according to this comment.
   
   
https://github.com/apache/beam/blob/92e3f353a85e11ddcdb9ae948511e7b8b662a762/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/CustomTimestampPolicyWithLimitedDelay.java#L75-L81
   
   Still, this change will break users with an intentional or unintentional 
dependency on implemented instead of designed behavior. :sweat_smile:
   
   Note that `ctx.getMessageBacklog()` may return 
`UnboundedReader.BACKLOG_UNKNOWN`.
   
   
https://github.com/apache/beam/blob/92e3f353a85e11ddcdb9ae948511e7b8b662a762/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaUnboundedReader.java#L513-L514
   
   The changes proposed in #39830 (port of `ReadFromKafkaDoFn` changes in 
#39285 to `KafkaUnboundedReader`) should make it less likely that the position 
gets ahead of the end offset (assuming that `currentLag()` is generally present 
after polling), because the end offset is no longer fetched by separate 
consumers and threads.
   
   I'm wondering if it makes sense for this policy to also advance the 
watermark to backlog check time (last succeeded backlog check time?) when 
`ctx.getMessageBacklog() <= 0` after the event time has advanced past 
`BoundedWindow.TIMESTAMP_MIN_VALUE`.



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