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]

Reply via email to