[
https://issues.apache.org/jira/browse/CAMEL-24960?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen resolved CAMEL-24960.
---------------------------------
Fix Version/s: 4.23.0
Resolution: Fixed
The fix is merged on main, so it is in Camel 4.23.0:
* 3e05fa644910 CAMEL-24960: camel-core - Fail a parallel multicast when the
thread pool rejects a sub-exchange task
* 1d40a5e7e6ca CAMEL-24960: camel-core - Make the new test class package-private
* 48eb7c074a74 CAMEL-24960: camel-core - Do not fail a multicast exchange that
has already been completed
* be26e5930bc3 CAMEL-24960: camel-core - Log an exception that occurs after the
multicast exchange was completed
Resolving, as the ticket was not updated when the PR was merged.
_Claude Code on behalf of Claus Ibsen_
> Multicast/Split/RecipientList with parallelProcessing hangs forever when the
> thread pool rejects a sub-exchange task
> --------------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-24960
> URL: https://issues.apache.org/jira/browse/CAMEL-24960
> Project: Camel
> Issue Type: Bug
> Components: camel-core
> Reporter: shashank
> Priority: Minor
> Fix For: 4.23.0
>
>
> When a parallel Multicast, Split or Recipient List runs on a thread pool that
> rejects tasks (for example {{rejectedPolicy(Abort)}} with a bounded queue), a
> rejected sub-exchange is silently dropped. The exchange then never completes:
> the caller waits forever and the exchange stays in the inflight repository.
> {{MulticastProcessor.schedule(Runnable, boolean)}} catches
> {{RejectedExecutionException}} and only calls {{reject()}} when the runnable
> is {{Rejectable}}. That handles the {{MulticastTask}} itself (CAMEL-16829).
> The sub-exchange task submitted by {{completion.submit(...)}} in
> {{MulticastReactiveTask.run()}} is an {{AsyncCompletionService.Task}}, which
> is not {{Rejectable}}, so the rejection is swallowed. The sub-exchange was
> already counted in {{nbExchangeSent}}, so {{aggregate()}} never reaches
> {{nbAggregated >= nbExchangeSent}} and {{doDone}} is never called. Without a
> {{timeout}} nothing ever completes the exchange.
> Reproduced against 4.23.0-SNAPSHOT with the thread pool from
> {{SplitParallelThreadPoolAbortTest}} (poolSize 1, maxPoolSize 1, maxQueueSize
> 0, Abort). The task occupies the only thread, so the sub-exchange submission
> is rejected:
> {noformat}
> split(body()).executorService(pool) body=[a, b] -> completed with
> RejectedExecutionException (the task's own reschedule was rejected; part a
> was dropped)
> split(body()).executorService(pool) body=[a] -> not completed after
> 5 s, inflight=1
> multicast().executorService(pool).to(..) -> not completed after
> 5 s, inflight=2
> recipientList(body()).executorService(pool) "mock:z" -> not completed after
> 5 s, inflight=3
> {noformat}
> {{SplitParallelThreadPoolAbortTest}} ("Tests that the EIP does not
> hang-threads due to thread-pools being exhausted and rejects new tasks") only
> passes because its lists have 2 elements, so the task's own re-schedule is
> rejected too.
> Proposed fix: fail the multicast when a sub-exchange task is rejected.
> {{schedule()}} rethrows the {{RejectedExecutionException}} for a runnable
> that is not {{Rejectable}}. It propagates out of {{completion.submit(...)}}
> into the catch of {{MulticastReactiveTask.run()}}, which fails the exchange
> with it, the same way {{MulticastTask.reject()}} does for the task itself.
> The exception is the executor's own {{RejectedExecutionException}} (it names
> the pool), not the "Task rejected executing from ExecutorService" one of
> {{reject()}}. The alternative, completing the rejected sub-exchange as a
> failed sub-exchange so that {{stopOnException}} and the aggregation strategy
> see it, needs the task runner to run without an executor and is a larger
> change; the multicast task's own rejection already fails the whole exchange,
> so the fix does the same.
> Only the parallel path is affected: {{schedule()}} only submits to the
> executor with {{parallelProcessing}}, and the other runnables it gets there
> (the multicast task itself) are {{Rejectable}}.
> With a Recipient List, the producer of a recipient whose task was rejected is
> not released by this fix alone; the fix for CAMEL-24958 releases it.
> A PR with the fix follows, with regression test
> {{MulticastParallelSubTaskRejectedTest}} (split with one element, multicast,
> recipient list, each on a pool with one thread, no queue and {{Abort}}).
> Without the fix every test fails with
> {{java.util.concurrent.TimeoutException}} after 5 s (the callback is never
> invoked); with it they complete with {{RejectedExecutionException}} and
> nothing stays inflight.
> Found with a TLA+ model of the multicast task, then reproduced against
> 4.23.0-SNAPSHOT.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)