[
https://issues.apache.org/jira/browse/CAMEL-24473?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen updated CAMEL-24473:
--------------------------------
Issue Type: Improvement (was: Bug)
> camel-kafka: consumer/producer readiness health check does not detect a
> downed broker under group.protocol=consumer (AsyncKafkaConsumer)
> ----------------------------------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-24473
> URL: https://issues.apache.org/jira/browse/CAMEL-24473
> Project: Camel
> Issue Type: Improvement
> Components: camel-kafka
> Reporter: Claus Ibsen
> Assignee: Omar Atie
> Priority: Major
> Fix For: 4.23.0
>
>
> h3. Problem
> The Kafka consumer readiness health check ({{KafkaConsumerHealthCheck}} ->
> {{KafkaFetchRecords.isReady()}}) determines broker connectivity via
> reflection into the Kafka client's internal {{ConsumerNetworkClient}} and
> calling {{hasReadyNodes(now)}}:
> {code:java}
> org.apache.kafka.clients.consumer.KafkaConsumer kc = ...;
> Object client =
> ReflectionHelper.getField(kc.getClass().getDeclaredField("delegate"), kc);
> ConsumerNetworkClient nc = (ConsumerNetworkClient) ReflectionHelper.getField(
> client.getClass().getDeclaredField("client"), client);
> ready = nc.hasReadyNodes(System.currentTimeMillis());
> {code}
> This works for the classic consumer protocol, where
> {{KafkaConsumer.delegate}} is a {{ClassicKafkaConsumer}} that has a {{private
> final ConsumerNetworkClient client}} field.
> Under the new consumer group protocol ({{group.protocol=consumer}}, KIP-848),
> {{KafkaConsumer.delegate}} is an {{AsyncKafkaConsumer}}, which has *no*
> {{client}} field and no {{ConsumerNetworkClient}} at all (networking runs on
> a background thread via {{NetworkClientDelegate}}). The
> {{getDeclaredField("client")}} call therefore throws
> {{NoSuchFieldException}}, which is caught and ignored, leaving the local
> {{ready}} at its default {{true}}.
> Result: with {{group.protocol=consumer}}, the consumer readiness health check
> reports UP even when the broker is down. It never transitions to DOWN.
> Verified against kafka-clients 4.3.1 (current {{<kafka-version>}} in
> parent/pom.xml):
> * {{ClassicKafkaConsumer}} -> has {{ConsumerNetworkClient client}}
> (reflection works)
> * {{AsyncKafkaConsumer}} -> no such field (reflection fails, swallowed, ready
> stays true)
> The producer path ({{KafkaProducer.isReady()}}) uses the same reflection
> pattern and should be reviewed for the equivalent gap.
> h3. Impact
> Deployments that enable the new Kafka consumer protocol
> ({{group.protocol=consumer}}) get a readiness probe that cannot detect a
> broker outage, defeating the purpose of the health check. Default protocol is
> still {{classic}} in Kafka 4.3.1, so the default configuration is unaffected
> today, but the new protocol is GA and adoption will grow.
> h3. Suggested direction
> Replace the brittle reflection with a protocol-agnostic readiness signal,
> e.g.:
> * use a public/stable Kafka client API where available, or
> * detect connectivity from poll/metadata state rather than reaching into
> private client internals, or
> * at minimum, handle the AsyncKafkaConsumer structure (background
> {{NetworkClientDelegate}}) explicitly and avoid defaulting {{ready=true}}
> when the expected internals are not found (fail-safe to a conservative state
> / log a warning).
> h3. Notes
> Discovered while investigating CAMEL-24466
> (KafkaConsumerHealthCheckIT.testReadinessWhenDown). That test failure had a
> different root cause (singleton test service whose shutdown() is a no-op),
> but the investigation surfaced this latent reflection gap.
> _Reported by Claude Code on behalf of davsclaus_
--
This message was sent by Atlassian Jira
(v8.20.10#820010)