[
https://issues.apache.org/jira/browse/CAMEL-25503?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Work on CAMEL-25503 started by shashank.
----------------------------------------
> camel-reactive-streams - Make the consumer Suspendable, so that a suspend
> does not stop it and wait for the drain
> -----------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-25503
> URL: https://issues.apache.org/jira/browse/CAMEL-25503
> Project: Camel
> Issue Type: Improvement
> Components: camel-reactive-streams
> Reporter: shashank
> Assignee: shashank
> Priority: Minor
>
> {{ReactiveStreamsConsumer}} (camel-reactive-streams, also used by
> camel-reactor and camel-rxjava) does not implement {{Suspendable}}. Camel
> therefore stops it whenever it would suspend it:
> * {{suspendRoute}} (route controller, JMX) stops the route: its status is
> {{Stopped}}, not {{Suspended}}, and a {{toStream}} request to it fails with
> {{No consumers attached to the stream}};
> * since CAMEL-25351 a stop waits until the items the consumer already took
> from the stream are routed (up to {{shutdownAwaitTermination}}, 10 s by
> default), so a suspend waits for them as well; items the publisher still
> sends for the outstanding demand are discarded with a WARN;
> * a route policy that suspends the consumer from one of its own exchanges
> ({{RoutePolicySupport.suspendOrStopConsumer}}, as
> {{ThrottlingInflightRoutePolicy}} does on the consumer's thread) stops it
> from its own thread pool: the queued items then fail with a
> {{RejectedExecutionException}} (one WARN each) and are lost.
> In the review of the CAMEL-25351 pull request
> (https://github.com/apache/camel/pull/27431), Claus Ibsen suggested a real
> suspend as a follow-up (detach and stop requesting, then resume), as the
> sjms, kafka and cxf consumers do (their suspend was fixed in CAMEL-25312,
> CAMEL-25313 and CAMEL-25314).
> h3. Proposed change
> * The consumer keeps the received items in a FIFO queue that its pool tasks
> take from in order (one task per item, as today), and implements
> {{Suspendable}}.
> * While it is suspending or suspended, the tasks leave the items queued and
> the subscriber requests nothing; items sent for the demand requested before
> the suspend are queued too (at most {{maxInflightExchanges}}). The exchanges
> being routed complete normally and the suspend waits for nothing.
> * {{doResume}} schedules the queued items and requests again; {{doStop}}
> schedules them before the drain added by CAMEL-25351, so a stop routes them;
> the stop is otherwise unchanged.
> Visible effects (upgrade guide and component doc): the route status after
> {{suspendRoute}} is {{Suspended}}; a graceful shutdown suspends the consumer
> before stopping it; a {{toStream}} request to a suspended route waits for the
> resume. With backpressure disabled ({{maxInflightExchanges}} not positive)
> the publisher cannot be paused, so items it sends while the consumer is
> suspended are kept in memory until the resume or stop.
> h3. Tests
> New {{ConsumerSuspendTest}}: a suspend does not wait for the queued items,
> holds them and requests nothing, and a stop routes them; route controller
> suspend/resume; {{toStream}} to a suspended route; a route policy suspending
> from the consumer's thread and resuming. Four of them fail on main (suspend
> timeout, {{expected: <Suspended> but was: <Stopped>}}, {{No consumers
> attached to the stream}}, consumer not suspended). camel-reactive-streams (75
> tests), camel-reactor (25) and camel-rxjava (25) pass with the change.
> Duplicate check (2026-10-08): JIRA text "reactive-streams suspend" (5 issues:
> CAMEL-25351 about the stop, the others unrelated), "ReactiveStreamsConsumer
> Suspendable" (CAMEL-25351 only), component camel-reactive-streams (13 issues,
> none about suspend). GitHub pull requests "reactive-streams suspend" and
> "ReactiveStreamsConsumer": #27431 and #27434 (both CAMEL-25351, stop only);
> no open pull request touches camel-reactive-streams.
> _Filed with Claude Code on behalf of allthingssecurity._
--
This message was sent by Atlassian Jira
(v8.20.10#820010)