[ 
https://issues.apache.org/jira/browse/CAMEL-25313?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Claus Ibsen resolved CAMEL-25313.
---------------------------------
    Resolution: Fixed

Merged in https://github.com/apache/camel/pull/27389 (commit 8723d9168286) for 
4.23.0.

_Claude Code on behalf of davsclaus_

> camel-kafka - a consumer resumed after a suspend may never consume again: the 
> fetcher thread ends while the consumer is suspended, and a pause or resume 
> request can be overwritten
> -----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25313
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25313
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-kafka
>            Reporter: shashank
>            Assignee: shashank
>            Priority: Major
>             Fix For: 4.23.0
>
>
> Suspending the route of a Kafka consumer ({{KafkaConsumer.doSuspend}}) only 
> sets the requested state of each {{KafkaFetchRecords}} task 
> ({{state.set(PAUSE_REQUESTED)}}, from the caller thread); resuming sets 
> {{RESUME_REQUESTED}}. The fetcher thread applies the request in 
> {{updateTaskState()}} after the next poll. Callers are the route controller 
> and JMX, and the route policies ({{ThrottlingExceptionRoutePolicy}} suspends 
> when the circuit opens and resumes from its half-open timer thread; 
> {{keepOpen}} suspends when the route starts).
> *1. The fetcher thread ends while the consumer is suspended.* {{run()}} 
> starts with {{if (!isKafkaConsumerRunnable()) return;}} and loops {{do { ... 
> startPolling(); } while ((canContinue || reconnect) && 
> isKafkaConsumerRunnable());}}, where {{isKafkaConsumerRunnable()}} is false 
> for a suspended consumer. The polling loop inside {{startPolling}} uses 
> {{isKafkaConsumerRunnableAndNotStopped()}} since CAMEL-18327, so it keeps 
> polling (paused) while suspended, but as soon as the thread leaves it, or 
> starts, while the consumer is suspended, it terminates ("Terminating 
> KafkaConsumer thread"). {{doResume}} then only sets the state of a task that 
> no thread runs: the route is started, the consumer never consumes again, and 
> nothing is logged as an error. This happens:
> * with {{ThrottlingExceptionRoutePolicy}} and {{keepOpen=true}}, always (not 
> a race) when the route starts with the CamelContext: the policy suspends the 
> consumer in {{onStart}}, and the fetcher tasks are only submitted once the 
> context has started ({{KafkaComponent.pendingConsumer}}), so {{run()}} 
> returns at once; setting {{keepOpen}} to false later (JMX) resumes a consumer 
> without a thread. This is the first part of CAMEL-19358 ("couldn't be resumed 
> after we changed the keepOpen to false");
> * when a poll fails while the consumer is suspended (default 
> {{pollOnError=ERROR_HANDLER}}: {{startPolling}} returns after handling the 
> exception), or with a reconnect ({{pollOnError=RECONNECT}}, or 
> {{breakOnFirstError}} when the suspend arrives during the batch).
> *2. A request made while the previous one is applied is lost.* 
> {{updateTaskState}} reads the state, calls {{consumer.pause(assignment)}} (or 
> seeks and calls {{consumer.resume(assignment)}}) and then sets {{PAUSED}} (or 
> {{RUNNING}}) unconditionally. A resume requested between the two is 
> overwritten: the Kafka consumer stays paused while the route is started 
> (never consumes again, {{isKafkaPaused}} stays true). A suspend requested 
> during the resume is overwritten by {{RUNNING}}: the consumer keeps consuming 
> while the route is suspended.
> h3. Reproduction
> {{KafkaConsumerSuspendResumeTest}} (camel-kafka, no broker: a 
> {{KafkaClientFactory}} returning a {{MockConsumer}} whose 
> {{pause}}/{{resume}} can be held on a latch to force the interleaving), on 
> main:
> {noformat}
> testStartedWithOpenCircuit (keepOpen=true at start, then keepOpen=false):
>   mock://result Received message count. Expected: <1> but was: <0>
> testPollErrorWhileSuspended (one poll throws while suspended, then resume):
>   mock://result Received message count. Expected: <1> but was: <0>
> testResumeWhileThePauseIsApplied (resumeRoute while the thread is in 
> consumer.pause()):
>   mock://result Received message count. Expected: <1> but was: <0>
> testSuspendWhileTheResumeIsApplied (suspendRoute while the thread is in 
> consumer.resume()):
>   The Kafka consumer must be paused while the route is suspended ==> 
> expected: <true> but was: <false> within 20 seconds.
> testStartedWithOpenCircuitConsumesNothing (a record on the topic when the 
> route starts with keepOpen=true):
>   ConditionTimeout: the consumer is never paused (the thread is gone)
> {noformat}
> {{testStartedWithOpenCircuitConsumesNothing}} also fails with only the first 
> three parts of the fix ({{Expected: <0> but was: <1>}}). No sleeps: latches, 
> Awaitility, {{MockEndpoint}}.
> The defect was found with a TLA+ model of the fetcher thread against route 
> suspend/resume ({{state.get()}}, the consumer call and {{state.set()}} as 
> separate steps; the run-loop conditions; a poll that can fail). "Once every 
> operation is done and no request is pending, a started route has a polling 
> thread with the consumer not paused, and a suspended route has its consumer 
> paused" is violated in 4 steps (suspend, the thread starts and returns, 
> resume), in 5 steps with a failing poll, and in 7 steps for the overwritten 
> resume (suspend, poll, read PAUSE_REQUESTED, resume, consumer.pause, set 
> PAUSED). With the fix it holds, as does the liveness property, also with two 
> suspend/resume cycles, failing polls and reconnects. The model has no 
> rebalance and abstracts the records ("alive and not paused"), so it does not 
> cover the rebalance-listener part, which the MockConsumer test does.
> h3. Proposed fix
> * {{updateTaskState}}: {{state.compareAndSet(PAUSE_REQUESTED, PAUSED)}} and 
> {{state.compareAndSet(RESUME_REQUESTED, RUNNING)}}, so a newer request is 
> handled on the next iteration.
> * {{run()}} and its do-while end only when the consumer is stopping or 
> stopped ({{isKafkaConsumerRunnableAndNotStopped()}}, as the polling loop): a 
> suspended consumer keeps its thread, which polls with the consumer paused.
> * After a (re)connect, {{state.compareAndSet(PAUSED, PAUSE_REQUESTED)}}: the 
> new Kafka consumer is not paused.
> * The rebalance listener ({{onPartitionsAssigned}}, on the fetcher thread 
> inside {{poll}}) pauses the assigned partitions when a pause is requested or 
> applied, so that the poll which assigns them does not return their records. 
> Without it a thread that starts or reconnects while suspended (now kept 
> alive) processes the first batch before {{updateTaskState}} pauses: with 
> {{keepOpen=true}} the records on the topic at startup would be consumed with 
> the circuit open (the other half of CAMEL-19358's report). Records of a batch 
> already polled when the suspend arrives are still processed first, as today.
> With the fix the new tests and the camel-kafka unit tests (224, 1 skipped) 
> pass; the integration tests need Docker and were not run.
> Affected: 4.14.x, 4.18.x and main (same code, GitHub contents API).
> Duplicate check (2026-10-04): JIRA "kafka" with suspend/resume and consumer 
> since 2022 (17 issues: CAMEL-19358 above, CAMEL-18327, CAMEL-18760, 
> CAMEL-18759, CAMEL-20227 offsets of the pausable consumer, none about the 
> thread or the lost requests), "KafkaFetchRecords" (CAMEL-25021 open, 
> oversized records; others unrelated). GitHub pull requests "kafka pause 
> resume": none open; open draft #26557 (exactly-once) does not change 
> {{KafkaFetchRecords}}.
> _Filed with Claude Code on behalf of allthingssecurity._



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to