Skip to content

fix: Race between run and cancellation in FrameProcessorExecutor. - #19006

Merged
gianm merged 6 commits into
apache:masterfrom
gianm:msq-bqfc-nicer
Sep 3, 2026
Merged

gianm merged 6 commits into
apache:masterfrom
gianm:msq-bqfc-nicer

Conversation

@gianm

@gianm gianm commented Feb 10, 2026 •

Copy link
Copy Markdown
Contributor

Previously, the ExecutorRunnable owned the processor only while running runProcessorNow(). This created races with cancellation. For example, if a processor was canceled and cleaned up while the ExecutorRunnable was calling "isFinished" or "readabilityFuture" on an input channel, it could lead to calling those operations on a closed channel.

This patch fixes it by having the ExecutorRunnable own the processor a bit longer, to cover the time that it may need to call methods on input and output channels. We now also also take care to not call methods on channels in the debug method logProcessorStatusString.

This patch also makes BlockingQueueFrameChannel more strict about closing, which helps ensure the above fix is working:

  1. Writable channel now rejects writes when closed, rather than when the reader has finished reading.

  2. Writable channel now rejects calls to all methods other than isClosed() when closed. Behavior changed in: writabilityFuture(), fail(), and close().

  3. Readable channel now rejects calls to all methods when closed. Behavior changed in: isFinished(), canRead(), read(), readabilityFuture(), and close().

@github-actions github-actions Bot added Area - Batch Ingestion Area - MSQ For multi stage queries - https://github.com/apache/druid/issues/12262 labels Feb 11, 2026
@gianm
gianm marked this pull request as draft March 3, 2026 22:21
Previously, the ExecutorRunnable owned the processor only while running
runProcessorNow(). This created races with cancellation. For example,
if a processor was canceled and cleaned up while the ExecutorRunnable
was calling "isFinished" or "readabilityFuture" on an input channel, it
could lead to calling those operations on a closed channel.

This patch fixes it by having the ExecutorRunnable own the processor
a bit longer, to cover the time that it may need to call methods on
input and output channels. We now also also take care to not call
methods on channels in the debug method logProcessorStatusString.

This patch also makes BlockingQueueFrameChannel more strict about
closing, which helps ensure the above fix is working:

1) Writable channel now rejects writes when closed, rather than when
   the reader has finished reading.

2) Writable channel now rejects calls to all methods other than
   isClosed() when closed. Behavior changed in: writabilityFuture(),
   fail(), and close().

3) Readable channel now rejects calls to all methods when closed.
   Behavior changed in: isFinished(), canRead(), read(),
   readabilityFuture(), and close().
@gianm gianm changed the title More defensive BlockingQueueFrameChannel. fix: Race between run and cancellation in FrameProcessorExecutor. Apr 2, 2026
@gianm

gianm commented Apr 2, 2026

Copy link
Copy Markdown
Contributor Author

The changes in this PR were originally scoped to just BlockingQueueFrameChannel, but they brought to light the race in FrameChannelExecutor, which is now the focus of the PR.

@gianm
gianm marked this pull request as ready for review April 2, 2026 19:56
throw DruidException.defensive("Closed, cannot call close() again");
}

readerClosed = true;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Readable close can leave writers permanently blocked

Readable.close() now only sets readerClosed and notifies the current writer future, but it leaves the queue contents in place and Writable.writabilityFuture() does not consider readerClosed. If the queue is full when the reader is closed, the notified writer can re-check writability, see the queue still full, create a new readyForWritingFuture, and wait forever because no reader will drain the queue or notify again. This can hang upstream processors or result producers when a consumer cancels or closes early. Consider clearing the queue as before, or making readerClosed close/fail the writable side so future waiters resume deterministically.

@gianm gianm Apr 26, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It's true although the old code is the same way. The old code cleared the queue here, but typically the queue is of length just 1 or 2, and it's easy for it to fill back up again. I believe isn't a problem in practice, because in production early-close scenarios (mainly LIMIT) processors are directly canceled by Worker#postCleanupStage. Generally we aren't relying on the channels themselves to propagate such information.

* Wait for all of the provided futures.
*/
public static <T> ReturnOrAwait<T> awaitAllFutures(final Collection<ListenableFuture<?>> futures)
public static <T> ReturnOrAwait<T> awaitAllFutures(final List<ListenableFuture<?>> futures)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Public awaitAllFutures signature breaks callers

Changing awaitAllFutures from Collection<ListenableFuture> to List> changes the public static method descriptor. Existing compiled extensions or downstream modules that call the Collection overload can fail with NoSuchMethodError, and source callers with another Collection type no longer compile. Keep the Collection overload, or retain the original signature and defensively copy to a List internally if list semantics are required.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This isn't meant to be a super-stable API, so it's OK.

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Follow-up handled. I rechecked the current code around both discussed threads and do not have additional inline replies or new high-confidence findings.


This is an automated review by Codex GPT-5


This is an automated review by Codex GPT-5

@github-actions

github-actions Bot commented Jul 1, 2026

Copy link
Copy Markdown

This pull request has been marked as stale due to 60 days of inactivity.
It will be closed in 4 weeks if no further activity occurs. If you think
that's incorrect or this pull request should instead be reviewed, please simply
write any comment. Even if closed, you can still revive the PR at any time or
discuss it on the dev@druid.apache.org list.
Thank you for your contributions.

@github-actions github-actions Bot added the stale label Jul 1, 2026
@github-actions

Copy link
Copy Markdown

This pull request/issue has been closed due to lack of activity. If you think that
is incorrect, or the pull request requires review, you can revive the PR at any time.

@github-actions github-actions Bot closed this Jul 30, 2026
@gianm gianm reopened this Jul 30, 2026
@gianm gianm removed the stale label Jul 30, 2026
@gianm
gianm merged commit ba84c74 into apache:master Sep 3, 2026
28 checks passed
@gianm
gianm deleted the msq-bqfc-nicer branch September 3, 2026 03:14
@github-actions github-actions Bot added this to the 39.0.0 milestone Sep 3, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Area - Batch Ingestion Area - MSQ For multi stage queries - https://github.com/apache/druid/issues/12262

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants