[
https://issues.apache.org/jira/browse/CAMEL-25029?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Federico Mariani updated CAMEL-25029:
-------------------------------------
Fix Version/s: 4.23.0
> 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
> Fix For: 4.23.0
>
>
> 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)