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]