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]