sjvanrossum commented on code in PR #40176:
URL: https://github.com/apache/beam/pull/40176#discussion_r4056770499
##########
sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/CustomTimestampPolicyWithLimitedDelay.java:
##########
@@ -59,6 +63,7 @@ public CustomTimestampPolicyWithLimitedDelay(
// 'previousWatermark' is not the same as maxEventTimestamp (e.g. it could
have been in future).
// Initialize it such that watermark before reading any event same as
previousWatermark.
maxEventTimestamp =
previousWatermark.orElse(BoundedWindow.TIMESTAMP_MIN_VALUE).plus(maxDelay);
+ lastWatermark = maxEventTimestamp.minus(maxDelay);
Review Comment:
```suggestion
lastWatermark =
previousWatermark.orElse(BoundedWindow.TIMESTAMP_MIN_VALUE);
maxEventTimestamp = lastWatermark.plus(maxDelay);
```
##########
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;
Review Comment:
Perhaps import `com.google.common.collect.Comparators` and use `max` since
`Instant` implements `Comparable`?
```suggestion
return (lastWatermark = Comparators.max(lastWatermark,
candidateWatermark(ctx, now)));
```
##########
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 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]