MSQ: Use interrupts for worker cancellation.#18095
Merged
Merged
Conversation
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.
clintropolis
approved these changes
Jun 26, 2025
clintropolis
left a comment
Member
There was a problem hiding this comment.
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 |
Member
There was a problem hiding this comment.
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 😅
Contributor
Author
There was a problem hiding this comment.
I had a similar train of thought when moving this code and reached the same conclusion.
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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:
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.
Move periodic controller status checking out of IndexerWorkerContext and into
a dedicated PeriodicControllerChecker class.
Remove stop(), awaitStop(), and controllerFailed() from the Worker interface.
Workers are now canceled via interruption only.
Remove controllerAlive tracking from WorkerImpl. It should no longer be necessary
now that workers are interrupted on controller failure.