kotwal-itpro opened a new pull request, #29397:
URL: https://github.com/apache/flink/pull/29397
## What is the purpose of the change
`AsyncWaitOperator` puts watermarks into its async queue so they are emitted
after the records that arrived before them, but it does not override
`processWatermarkStatus`. `AbstractStreamOperator#processWatermarkStatus`
forwards the status right away, so a `WatermarkStatus` can overtake watermarks
that are still queued. For the input `[record, W100, IDLE]` the operator emits
`[IDLE, record, W100]`.
A `StatusWatermarkValve` downstream ignores watermarks that arrive on an
idle channel, so W100 is lost, not just delayed. If it was the last watermark
on that channel (e.g. `MAX_WATERMARK` at drain), downstream event time never
advances past it.
## Brief change log
- Added `WatermarkStatusQueueEntry`, a queue entry that is always
completed and emits a `WatermarkStatus`.
- `OrderedStreamElementQueue` and `UnorderedStreamElementQueue` accept
watermark statuses. The unordered queue gives a status its own segment, like a
watermark, so records cannot be reordered across it.
- `AsyncWaitOperator#processWatermarkStatus` adds the status to the queue
and emits completed elements, the same way as `processWatermark`. A status with
an empty queue is still emitted immediately.
- Statuses that are still queued at checkpoint time are part of the
snapshotted queue elements, and `open()` replays them on restore.
## Verifying this change
This change added tests and can be verified as follows:
-
`AsyncWaitOperatorTest#testWatermarkStatusDoesNotOvertakeWatermark{Ordered,Unordered}`:
`[record, W100, IDLE]` with the async result withheld emits nothing until the
record completes, then `[record, W100, IDLE]`.
- `AsyncWaitOperatorTest#testWatermarkStatusIsRestoredInOrder`: a status
that is queued at snapshot time is restored and emitted after the restored
record and watermark.
-
`AsyncWaitOperatorTest#testWatermarkStatusWithEmptyQueueIsEmittedImmediately`:
no extra latency when nothing is queued.
- `StreamElementQueueTest#testWatermarkStatusIsEmittedInOrder` (both
queues) and
`UnorderedStreamElementQueueTest#testWatermarkStatusCompletionOrder`.
All new tests except the empty-queue one fail without the change.
`AsyncWaitOperatorTest`, `StreamElementQueueTest`,
`OrderedStreamElementQueueTest` and `UnorderedStreamElementQueueTest` pass.
## Does this pull request potentially affect one of the following parts:
- Dependencies (does it add or upgrade a dependency): no
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: no
- The serializers: no (the operator state already uses
`StreamElementSerializer`, which supports `WatermarkStatus`)
- The runtime per-record code paths (performance sensitive): no (only the
watermark status path)
- Anything that affects deployment or recovery: JobManager (and its
components), Checkpointing, Kubernetes/Yarn, ZooKeeper: yes. The async operator
state can now contain `WatermarkStatus` elements. Restoring such a checkpoint
with a version without this change fails in `AsyncWaitOperator#open()` with
"Unknown record type".
- The S3 file system connector: no
## Documentation
- Does this pull request introduce a new feature? no
- If yes, how is the feature documented? not applicable
`RowTimeMiniBatchAssignerOperator` has the same shape (it buffers watermarks
but forwards the status immediately), as noted in the JIRA. I kept this PR to
`AsyncWaitOperator` and can follow up on that one separately.
---
##### Was generative AI tooling used to co-author this PR?
- [X] Yes (please specify the tool below)
Generated-by: Claude Opus 5.5
--
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]