shashank created CAMEL-25503:
--------------------------------
Summary: 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
{{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)