[
https://issues.apache.org/jira/browse/CAMEL-25313?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen reassigned CAMEL-25313:
-----------------------------------
Assignee: shashank
> 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)