davsclaus commented on code in PR #27002:
URL: https://github.com/apache/camel/pull/27002#discussion_r4130733431
##########
core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java:
##########
@@ -188,27 +197,65 @@ public boolean process(Exchange exchange, AsyncCallback
callback) {
}
/**
- * Submits the onCompletion task to the thread pool (parallel processing).
The task is counted as pending from when
- * it is submitted until it is done, so a graceful shutdown waits for it.
+ * Submits the onCompletion task of the given copy to the thread pool
(parallel processing). The task is counted as
+ * pending from when it is submitted until it is done, so a graceful
shutdown waits for it.
+ * <p>
+ * The copy may hold its own reference to a stream cache (see {@link
#prepareExchange(Exchange)}), which is released
+ * when the copy is done. When the task never runs (the thread pool
rejects it, discards it because it is shut down,
+ * or drops it when the processor shuts its thread pool down), it is no
longer counted as pending and the reference
+ * is released instead, as otherwise a spooled file would be kept until
the stream caching strategy is stopped.
*/
@SuppressWarnings("deprecation")
- private void submitTask(Runnable task) {
+ private void submitTask(Exchange copy, Runnable task) {
+ ParallelTask parallelTask = new ParallelTask(copy, task);
taskCount.increment();
- Runnable counted = () -> {
- try {
- task.run();
- } finally {
- taskCount.decrement();
- }
- };
+ pendingTasks.add(parallelTask);
try {
// Deprecated since 4.19.0
- executorService.submit(prepareMDCParallelTask(camelContext,
counted));
+ executorService.submit(prepareMDCParallelTask(camelContext,
parallelTask));
} catch (RuntimeException e) {
// the task will not run
- taskCount.decrement();
+ parallelTask.discard();
throw e;
}
+ if (executorService.isShutdown()) {
Review Comment:
Minor, non-blocking: if `submit` queued the task on a running pool and
another thread then calls a graceful `shutdown()` before this check, the task
is discarded here. But `ThreadPoolExecutor.shutdown()` still runs queued tasks,
so when the worker later calls `run()` the task has already been removed and
that onCompletion is silently lost. The window is tiny and only during
shutdown. Could the check be narrowed (e.g. for a `ThreadPoolExecutor`, discard
only when the task is no longer in `getQueue()`), or the trade-off documented?
--
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]