1996fanrui opened a new pull request, #29056:
URL: https://github.com/apache/flink/pull/29056
## What is the purpose of the change
A blocked `PipelinedSubpartition` still hands out priority buffers
(unaligned-checkpoint barriers) in
`pollBuffer()`, but `isDataAvailableUnsafe()` and
`getBuffersInBacklogUnsafe()` still gate on `!isBlocked`.
A credited remote reader decides availability through those methods, so its
one-shot priority notification
is evaluated as "not available", the reader is not enqueued, and the barrier
is not delivered until recovery
ends (`resumeConsumption()`) — the checkpoint can time out. The local
channel path is already handled via
`hasPendingPriorityEvent`; the remote path was not.
## Brief change log
- [FLINK-40519] Extract `PipelinedSubpartition#isBlockedForDelivery()` and
use it in `pollBuffer`, `isDataAvailableUnsafe` and `getBuffersInBacklogUnsafe`
so a queued priority barrier is delivered to a credited reader while the
subpartition is blocked.
## Verifying this change
This change added a unit test:
- Added a `PipelinedSubpartitionTest` case asserting that, while blocked
with a priority barrier queued, the credited-reader availability check reports
available and the barrier is pollable.
## 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 runtime per-record code paths (performance sensitive): no
- Anything that affects deployment or recovery: yes (unaligned-checkpoint
barrier delivery during recovery)
- 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
---
##### Was generative AI tooling used to co-author this PR?
- [ ] Yes (please specify the tool below)
--
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]