davsclaus commented on code in PR #27636:
URL: https://github.com/apache/camel/pull/27636#discussion_r4236604030
##########
components/camel-reactive-streams/src/main/java/org/apache/camel/component/reactive/streams/ReactiveStreamsConsumer.java:
##########
@@ -61,17 +75,47 @@ protected void doStart() throws Exception {
getEndpoint().getEndpointUri(), poolSize);
}
- this.service.attachCamelConsumer(endpoint.getStream(), this);
+ // the items left queued by a stop that timed out are routed now
+ scheduleQueuedItems();
+
+ this.subscriber =
this.service.attachCamelConsumer(endpoint.getStream(), this);
+ }
+
+ @Override
+ protected void doSuspend() throws Exception {
Review Comment:
In a graceful *shutdown* (not a suspend-only), this suspend changes the
point at which the queued items get routed. Here is how
`DefaultShutdownStrategy.ShutdownTask` handles it:
1. First pass, routes in reverse startup order: this consumer is
`Suspendable`, so it is suspended and put on the deferred list. A downstream
`DirectConsumer`/`SedaConsumer` returns `deferShutdown=true` and is also
deferred, but it keeps running.
2. Inflight wait: only the inflight repository and
`ShutdownAware.getPendingExchangesSize` are counted. The items held in `queued`
are in neither, so the wait ends right away.
3. Deferred consumers are stopped in list order. For
`from("reactive-streams:in").to("direct:sub")` defined before
`from("direct:sub")`, `direct:sub` comes first and is removed. Only then does
this `doStop()` call `scheduleQueuedItems()`, and those items hit a
`direct:sub` that has no consumer. With `block=true`,
`DirectComponent.getConsumer` waits up to 30 s per item, `shutdownGraceful`
gives up after `shutdownAwaitTermination`, and the items are lost.
On main the consumer is stopped in pass 1, so the drain in `doStop` runs
while `direct:sub` is still up.
One way to keep the suspend for route controller and route policy use, but
still drain during a shutdown: implement `ShutdownAware` the way `SedaConsumer`
does. `deferShutdown` returns `false`, so the strategy still suspends. In
`getPendingExchangesSize(boolean suspendOnly)`, when `!suspendOnly`, switch
into a "draining" mode: `routeQueuedItem` routes even while suspended,
`refill()` still requests nothing, and you call `scheduleQueuedItems()` and
return `queued.size()` plus the items being routed. When `suspendOnly`, return
0 as now. The shutdown wait then covers the queued items while the deferred
downstream consumers are still running, and the shutdown strategy's timeout
applies to them as well.
--
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]