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


##########
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:
   Unconditionally changing this may break existing users.
   If a new partition is added or an existing partition is cleared before 
running a pipeline with the intent being to produce records to the partition 
after the pipeline is running and healthy, then there's a good reason to hold 
the watermark at `TIMESTAMP_MIN_VALUE`.
   
   I'd consider making the proposed change configurable or splitting it out to 
a separate class instead of changing it unconditionally.



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