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

Work on CAMEL-25029 started by Federico Mariani.
------------------------------------------------
> camel-kafka - Add a consumer for Kafka share groups (KIP-932, Queues for 
> Kafka)
> -------------------------------------------------------------------------------
>
>                 Key: CAMEL-25029
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25029
>             Project: Camel
>          Issue Type: New Feature
>          Components: camel-kafka
>            Reporter: Federico Mariani
>            Assignee: Federico Mariani
>            Priority: Minor
>
> h3. Motivation
> Kafka share groups (KIP-932, "Queues for Kafka") let consumers work through 
> the records of a topic like a queue:
> * several consumers can read the same partition, so the number of consumers 
> is no longer capped by the number of partitions
> * the application acknowledges each record: {{ACCEPT}}, {{RELEASE}} 
> (redeliver), {{REJECT}} (discard) or {{RENEW}} (extend the acquisition lock)
> * the broker handles redelivery and counts deliveries (exposed as 
> {{ConsumerRecord.deliveryCount()}}); a failing record does not block its 
> partition, and records are archived after 
> {{group.share.delivery.count.limit}} deliveries
> * no offsets, no seeking, no topic patterns, no ordering guarantees, no 
> transactions
> The client API is available in the kafka-clients version Camel already uses 
> (4.3.1: {{KafkaShareConsumer}}, {{AcknowledgeType}}), so this needs no 
> dependency upgrade.
> Typical uses:
> * work queues for slow, independent tasks (LLM/embedding calls, OCR, 
> rate-limited APIs) that need to scale beyond the partition count
> * replacing a JMS/AMQP work queue with Kafka, without a second broker
> * consuming the same topic in two ways: an ordered consumer group for 
> projections, and a share group of workers for unordered side effects
> h3. Proposal
> A new consumer-only scheme, {{kafka-share:topic}}, in the camel-kafka module, 
> with its own endpoint, configuration and consumer.
> It would *not* be an option on {{kafka:}}, because about 14 of the 34 
> existing consumer options do not apply to a share consumer ({{seekTo}}, 
> {{autoOffsetReset}}, {{autoCommitEnable}}, {{autoCommitIntervalMs}}, 
> {{allowManualCommit}}, {{commitTimeoutMs}}, {{partitionAssignor}}, 
> {{groupProtocol}}, {{groupRemoteAssignor}}, {{groupInstanceId}}, 
> {{topicIsPattern}}, {{isolationLevel}}, {{batching}}, 
> {{batchingIntervalMs}}). The current consumer is built around offsets, commit 
> managers and rebalance listeners, none of which exist here. A separate scheme 
> keeps the catalog and tooling accurate. Brokers, security, serdes, 
> {{KafkaClientFactory}}, the header filter strategy and the health check 
> infrastructure would be shared within the module. Producing stays on 
> {{kafka:}}.
> Acknowledgement is derived from the exchange outcome (explicit 
> acknowledgement mode):
> || Exchange outcome || Acknowledgement ||
> | completed, or failed with {{handled(true)}} | {{ACCEPT}} |
> | failed, or rolled back (transacted route) | configurable, default 
> {{RELEASE}} |
> | header set by the route | {{ACCEPT}} / {{RELEASE}} / {{REJECT}} |
> | still inflight when the lock is about to expire (optional) | {{RENEW}} |
> The default {{RELEASE}} lets the broker redeliver, possibly to another 
> instance; the broker's delivery count limit prevents endless redelivery of 
> poison messages.
> Headers: delivery count, topic, partition, offset, key, timestamp.
> h3. Examples
> _Option and header names are illustrative._
> *1. Work queue that scales past the partition count*
> {code:java}
> from("kafka-share:support-tickets?brokers={{kafka.brokers}}&groupId=ticket-triage&consumersCount=20")
>     .to("langchain4j-chat:triage")
>     .to("kafka:tickets-triaged?brokers={{kafka.brokers}}");
> {code}
> 20 concurrent consumers on a topic with 3 partitions; a slow ticket does not 
> hold up the others.
> *2. Broker-driven redelivery with a delivery-count cutoff*
> {code:java}
> errorHandler(noErrorHandler()); // no local retries: the broker redelivers, 
> possibly to another pod
> onException(InvalidInvoiceException.class)
>     .handled(true)                                        // handled -> ACCEPT
>     .to("kafka:invoices-invalid?brokers={{kafka.brokers}}");
> from("kafka-share:invoices?brokers={{kafka.brokers}}&groupId=invoice-ocr")
>     
> .filter(header(KafkaConstants.SHARE_DELIVERY_COUNT).isGreaterThanOrEqualTo(4))
>         .setHeader(KafkaConstants.SHARE_ACKNOWLEDGE, constant("REJECT"))
>         .to("kafka:invoices-dlq?brokers={{kafka.brokers}}")
>         .stop()
>     .end()
>     .to("http://ocr-service/extract";)
>     .to("kafka:invoices-extracted?brokers={{kafka.brokers}}");
> {code}
> h3. Notes and follow-ups
> * *Documentation - {{consumersCount}}:* each consumer thread owns one 
> {{KafkaShareConsumer}} (the client is not thread-safe). Unlike classic 
> consumers, consumers beyond the partition count are not idle: all of them 
> receive records.
> * *Documentation - delivery semantics:* delivery is at-least-once. A record 
> can be processed again if its acquisition lock expires during processing, if 
> the JVM stops after the route's side effects but before the acknowledgement 
> reaches the broker, if committing the acknowledgements fails, or if the route 
> releases a record after a partial side effect. Pair with the Idempotent 
> Consumer where duplicates matter.
> * *Manual acknowledgement:* should use the planned acknowledgement API, the 
> share-group equivalent of {{KafkaManualCommit}}. Acknowledgements made from 
> another thread (async hand-off) have to be applied on the polling thread.
> * *Follow-up - virtual threads:* split polling from processing, i.e. a few 
> pollers, each record of a poll processed on a virtual thread, 
> acknowledgements queued back to the poller. This gives high concurrency for 
> I/O-bound routes without one Kafka client per concurrent exchange.
> * Health check criteria, since there is no partition assignment to report.
> h3. References
> * KIP-932: 
> https://cwiki.apache.org/confluence/display/KAFKA/KIP-932%3A+Queues+for+Kafka



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

Reply via email to