Skip to content

MSQ: Use interrupts for worker cancellation.#18095

Merged
gianm merged 2 commits into
apache:masterfrom
gianm:msq-stop-interrupt
Jun 27, 2025
Merged

MSQ: Use interrupts for worker cancellation.#18095
gianm merged 2 commits into
apache:masterfrom
gianm:msq-stop-interrupt

Conversation

@gianm

@gianm gianm commented Jun 8, 2025

Copy link
Copy Markdown
Contributor

This patch simplifies the worker lifecycle and focuses on using interruption instead of various flags to manage cancellation. This should enable workers to cancel themselves more quickly, and is particularly important when they may be waiting for a resource to become available (such as merge buffers as in #18075).

Changes:

  1. Introduce WorkerRunRef to manage worker execution and cancellation. This class
    encapsulates the logic for running a worker in an executor and canceling it
    via thread interruption.

  2. Move periodic controller status checking out of IndexerWorkerContext and into
    a dedicated PeriodicControllerChecker class.

  3. Remove stop(), awaitStop(), and controllerFailed() from the Worker interface.
    Workers are now canceled via interruption only.

  4. Remove controllerAlive tracking from WorkerImpl. It should no longer be necessary
    now that workers are interrupted on controller failure.

This patch simplifies the worker lifecycle and focuses on using interruption
instead of various flags to manage cancellation. This should enable workers to
cancel themselves more quickly, and is particularly important when they may
be waiting for a resource to become available (such as merge buffers as in apache#18075).

Changes:

1) Introduce WorkerRunRef to manage worker execution and cancellation. This class
   encapsulates the logic for running a worker in an executor and canceling it
   via thread interruption.

2) Move periodic controller status checking out of IndexerWorkerContext and into
   a dedicated PeriodicControllerChecker class.

3) Remove stop(), awaitStop(), and controllerFailed() from the Worker interface.
   Workers are now canceled via interruption only.

4) Remove controllerAlive tracking from WorkerImpl. It should no longer be necessary
   now that workers are interrupted on controller failure.
@github-actions github-actions Bot added Area - Batch Ingestion Area - MSQ For multi stage queries - https://github.com/apache/druid/issues/12262 labels Jun 8, 2025

@clintropolis clintropolis 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.

lgtm, this seems to tidy some stuff up nicely too 👍

import java.util.concurrent.ExecutorService;
import java.util.concurrent.ThreadLocalRandom;

public class PeriodicControllerChecker implements Closeable

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.

naming nit: based on name I would kind of expect this to use a scheduled executor instead of a while loop/sleep. however, this does run at a fixed frequency so i don't feel very strongly either way (and this makes more sense than using some kind of scheduled executor too, so not suggesting changing anything). naming is the worst 😅

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.

I had a similar train of thought when moving this code and reached the same conclusion.

@gianm
gianm merged commit 8e39b40 into apache:master Jun 27, 2025
42 checks passed
@gianm
gianm deleted the msq-stop-interrupt branch June 27, 2025 15:25
@capistrant capistrant added this to the 34.0.0 milestone Jul 22, 2025
gianm added a commit to gianm/druid that referenced this pull request Jan 20, 2026
Patch apache#18095 replaced various previously-existing worker cancellation
mechanisms with an interrupt of the worker thread. Unfortunately, in
debugging some stuck tests, it has been revealed that there is at least
one scenario where an InterruptedException could be swallowed when the
main worker thread handles a new work order. This defeated the
cancellation mechanism, leading the worker to run forever.

This patch does three things to improve robustness:

1) Fix the specific code path in RunAllFullyWidget that was swallowing
   InterruptedException in the problematic test.

2) Add a non-interrupt-based cancellation mechanism: a stop() call that
   throws an exception into the main thread's kernel manipulation queue.
   This is useful as a failsafe in case the interrupt gets lost in some
   other code path that may not have been discovered yet.

3) Update MSQTestBase to wait for workers to exit before moving on to
   the next test, and fail if they don't exit within 10 seconds. This is
   preferable to the tests running forever and possibly polluting the
   shared executor.
gianm added a commit to gianm/druid that referenced this pull request Jan 20, 2026
Patch apache#18095 replaced various previously-existing worker cancellation
mechanisms with an interrupt of the worker thread. Unfortunately, in
debugging some stuck tests, it has been revealed that there is at least
one scenario where an InterruptedException could be swallowed when the
main worker thread handles a new work order. This defeated the
cancellation mechanism, leading the worker to run forever.

This patch does three things to improve robustness:

1) Fix the specific code path in RunAllFullyWidget that was swallowing
   InterruptedException in the problematic test.

2) Add a non-interrupt-based cancellation mechanism: a stop() call that
   throws an exception into the main thread's kernel manipulation queue.
   This is useful as a failsafe in case the interrupt gets lost in some
   other code path that may not have been discovered yet.

3) Update MSQTestBase to wait for workers to exit before moving on to
   the next test, and fail if they don't exit within 10 seconds. This is
   preferable to the tests running forever and possibly polluting the
   shared executor.
gianm added a commit that referenced this pull request Jan 20, 2026
Patch #18095 replaced various previously-existing worker cancellation
mechanisms with an interrupt of the worker thread. Unfortunately, in
debugging some stuck tests, it has been revealed that there is at least
one scenario where an InterruptedException could be swallowed when the
main worker thread handles a new work order. This defeated the
cancellation mechanism, leading the worker to run forever.

This patch does three things to improve robustness:

1) Fix the specific code path in RunAllFullyWidget that was swallowing
   InterruptedException in the problematic test.

2) Add a non-interrupt-based cancellation mechanism: a stop() call that
   throws an exception into the main thread's kernel manipulation queue.
   This is useful as a failsafe in case the interrupt gets lost in some
   other code path that may not have been discovered yet.

3) Update MSQTestBase to wait for workers to exit before moving on to
   the next test, and fail if they don't exit within 10 seconds. This is
   preferable to the tests running forever and possibly polluting the
   shared executor.
gianm added a commit to gianm/druid that referenced this pull request Mar 30, 2026
In PRs apache#18095 and apache#18931, worker cancellation was switched to use
interrupts with a lightweight non-interrupt-based failsafe. This patch
implements a similar idea for controllers, to aid in more prompt
cancellation in cases where the controller is blocking on something.

The main change is to track the controller thread in ControllerHolder
and interrupt it on cancel(), in addition to calling controller.stop().
ControllerHolder is also moved from dart.controller to msq.exec, since
it is now a shared class, no longer Dart-specific. In addition, the
"workerOffline" logic is moved to Dart's ControllerMessageListener,
since that really is Dart-specific.
gianm added a commit to gianm/druid that referenced this pull request Mar 30, 2026
In PRs apache#18095 and apache#18931, worker cancellation was switched to use
interrupts with a lightweight non-interrupt-based failsafe. This patch
implements a similar idea for controllers, to aid in more prompt
cancellation in cases where the controller is blocking on something.

The main change is to track the controller thread in ControllerHolder
and interrupt it on cancel(), in addition to calling controller.stop().
ControllerHolder is also moved from dart.controller to msq.exec, since
it is now a shared class, no longer Dart-specific. In addition, the
"workerOffline" logic is moved to Dart's ControllerMessageListener,
since that really is Dart-specific.
gianm added a commit that referenced this pull request Apr 3, 2026
In PRs #18095 and #18931, worker cancellation was switched to use
interrupts with a lightweight non-interrupt-based failsafe. This patch
implements a similar idea for controllers, to aid in more prompt
cancellation in cases where the controller is blocking on something.

The main change is to track the controller thread in ControllerHolder
and interrupt it on cancel(), in addition to calling controller.stop().
ControllerHolder is also moved from dart.controller to msq.exec, since
it is now a shared class, no longer Dart-specific. In addition, the
"workerOffline" logic is moved to Dart's ControllerMessageListener,
since that really is Dart-specific.
riovic918data pushed a commit to riovic918data/druid that referenced this pull request Jun 12, 2026
* MSQ: Use interrupts for worker cancellation.

This patch simplifies the worker lifecycle and focuses on using interruption
instead of various flags to manage cancellation. This should enable workers to
cancel themselves more quickly, and is particularly important when they may
be waiting for a resource to become available (such as merge buffers as in apache#18075).

Changes:

1) Introduce WorkerRunRef to manage worker execution and cancellation. This class
   encapsulates the logic for running a worker in an executor and canceling it
   via thread interruption.

2) Move periodic controller status checking out of IndexerWorkerContext and into
   a dedicated PeriodicControllerChecker class.

3) Remove stop(), awaitStop(), and controllerFailed() from the Worker interface.
   Workers are now canceled via interruption only.

4) Remove controllerAlive tracking from WorkerImpl. It should no longer be necessary
   now that workers are interrupted on controller failure.

* Style fixes.
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