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


##########
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:
   Hi @sjvanrossum , thanks for taking a look 🙇‍♂️  .
   
   Pushed an update. It's opt-in constructor flag, default off, so nothing
   changes for existing users. That should cover the "might break users" worry.
   
   I also dropped the monotonicity clamp, which is where both your code 
suggestions
   were, sorry about that. The reason: the watermark can already go backwards 
today.
   Advance while idle, then a record arrives and you fall back to
   `maxEventTimestamp - maxDelay`, which can be well behind the idle value. 
That's
   existing behaviour with the gate in place, so fixing it felt like a separate
   thing. The diff is now just the gate condition plus a constructor overload.
   
   On `BACKLOG_UNKNOWN`: -1 fails `== 0`, so the watermark just holds. Safe, and
   this PR doesn't change it.
   
   On advancing at `<= 0` — I'd keep that as a separate issue, since
   `BACKLOG_UNKNOWN` means "I don't know if I'm caught up", which isn't the 
same as
   `== 0`.
   
   However, right now the flag is only reachable by building the policy 
yourself and
   passing it through `withTimestampPolicyFactory`. I left it off from 
`KafkaIO.Read` to
   keep this special behaviour away from default use cases. 



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