[ 
https://issues.apache.org/jira/browse/CAMEL-25029?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18124365#comment-18124365
 ] 

Federico Mariani commented on CAMEL-25029:
------------------------------------------

h3. Implementation strategy and sequencing

After reviewing the module layout question (whether to split {{camel-kafka}} 
into a common module plus two components, or to add {{kafka-share}} as a second 
scheme inside {{camel-kafka}}), this is the plan.

h4. Why the share consumer cannot reuse the classic consumer classes

* {{ShareConsumerConfig}} (kafka-clients 4.3.1) throws a {{ConfigException}} 
when any of these keys is present: {{auto.offset.reset}}, 
{{enable.auto.commit}}, {{group.instance.id}}, {{isolation.level}}, 
{{partition.assignment.strategy}}, {{interceptor.classes}}, 
{{session.timeout.ms}}, {{heartbeat.interval.ms}}, {{group.protocol}}, 
{{group.remote.assignor}}. {{KafkaConfiguration.createConsumerProperties()}} 
sets six of them unconditionally, so the share consumer needs its own 
configuration class that shares only the common block (brokers, client id, SSL, 
SASL, serdes, header filter, backoffs, additional properties).
* {{ShareConsumer}} has no {{pause}}/{{resume}}, {{seek}}, {{assignment}}, 
{{committed}}, topic patterns, rebalance listener or transactions, so the fetch 
loop, commit managers, manual commit, pausable consumer, resume strategy and 
dev console {{committed}} view do not apply.
* Acknowledgement mode will be {{share.acknowledgement.mode=explicit}}; every 
record of a poll is acknowledged on the polling thread before the next poll. 
{{MockShareConsumer}} ships in kafka-clients and will be used for unit tests.

h4. Precedents in the code base

* {{camel-infinispan}} split into {{camel-infinispan-common}}, 
{{camel-infinispan}} and {{camel-infinispan-embedded}} (CAMEL-12489).
* {{camel-ftp-common}} extracted from {{camel-ftp}} (CAMEL-23164), package 
unchanged, upgrade guide note only.
* Two schemes in one module: {{google-mail}} / {{google-mail-stream}}, 
{{aws2-kinesis}} / {{aws2-kinesis-firehose}}, {{ftp}} / {{ftps}} / {{sftp}}.
* The catalog tooling already ignores {{-common}} directories and the Spring 
Boot generator creates no starter for them, so a flat {{camel-kafka-common}} 
needs no tooling registration.

h4. Sequencing: three PRs on {{main}} for 4.23.0

# *Extraction inside {{camel-kafka}}* (refactor only, no behaviour or option 
change, no upgrade guide entry):
#* {{KafkaClientConfiguration}} (abstract, {{@UriParams}}): the ~45 common 
options and the SSL/SASL/additional-properties builders moved out of 
{{KafkaConfiguration}}, which extends it. Field order is preserved so 
{{kafka.json}} and the doc table stay identical.
#* {{AbstractKafkaComponent}}: client factory, poll exception strategy, 
create/subscribe backoff options, global SSL, startup deferral of consumers.
#* {{KafkaClientFactory.getShareConsumer(Properties)}} default method; 
{{KafkaSecurityConfigurer}} and {{KafkaRecordProcessor.propagateHeaders}} take 
the new base type.
#* Create/subscribe reconnect backoff and health state extracted from 
{{KafkaFetchRecords}} into a helper; the pause/resume state machine 
(CAMEL-25313) stays in {{KafkaFetchRecords}}.
#* Small interface for the error strategies that currently take 
{{KafkaFetchRecords}}.
# *Move the shared layer to {{camel-kafka-common}}* (flat layout next to 
{{camel-kafka}}, same packages, pure {{git mv}}): common pom takes over 
{{kafka-clients}} and the lz4 replacement, attaches a {{test-jar}} with the 
shared IT support; {{camel-kafka}} depends on it. Wiring: 
{{components/pom.xml}}, {{parent/pom.xml}}, regenerated BOM, allcomponents and 
sbom, labeler globs, an upgrade guide entry worded like the 4.19 
{{camel-ftp-common}} one. The decision whether this PR is wanted at all can be 
taken after PR 1: if the project prefers a single artifact, {{kafka-share}} 
lands inside {{camel-kafka}} as a second scheme and PR 1 and PR 3 are unchanged.
# *{{camel-kafka-share}} component* ({{kafka-share:topic}}, consumer only, 
package {{org.apache.camel.component.kafka.share}}), built only on 
{{camel-kafka-common}}:
#* {{KafkaShareConfiguration extends KafkaClientConfiguration}} with: {{topic}} 
(list, no pattern), {{groupId}} (required, no random default), 
{{consumersCount}}, {{pollTimeoutMs}}, {{maxPollRecords}}, fetch options, 
{{onFailure}} (ACCEPT/RELEASE/REJECT, default RELEASE), {{onRollback}} (default 
RELEASE), {{acknowledgementCommit}} (sync/async/poll), {{commitTimeoutMs}}, 
{{pollOnError}}, {{allowManualAcknowledgement}}. Never emits the ten rejected 
keys (unit test).
#* Acknowledgement derived from the exchange outcome: header 
{{CamelKafkaShareAcknowledge}} wins, then the manual acknowledgement object, 
then completed or handled = ACCEPT, failed = {{onFailure}}, rolled back = 
{{onRollback}}, in flight at stop = RELEASE.
#* Headers: {{CamelKafkaShareDeliveryCount}} (from 
{{ConsumerRecord.deliveryCount()}}), plus the existing topic, partition, 
offset, key, timestamp and headers constants.
#* {{KafkaShareConsumer}} with one {{KafkaShareFetchRecords}} per 
{{consumersCount}} (reconnect helper from PR 1), suspend = stop polling and 
release in-flight records, health check (connected and subscribed, no 
assignment), {{kafka-share}} dev console with delivery and acknowledgement 
counters.
#* Producer side stays on {{kafka:}}; {{createProducer}} throws.
#* Tests: unit with {{MockShareConsumer}} (acknowledgement table, header 
override, manual API, release on stop, per-partition commit errors), ITs on the 
existing 4.3.1 test-infra broker (N consumers on one partition all receive, 
delivery count on RELEASE, REJECT not redelivered, health check bad port and 
recovery, stop releases to a second consumer, SASL path). Redpanda/Strimzi 
flavours skipped.
#* Docs: new {{kafka-share-component.adoc}} (when to use, broker requirements, 
options, headers, acknowledgement table, at-least-once semantics, 
{{consumersCount}} semantics, differences from {{kafka:}}), pointer paragraph 
in the {{kafka}} page.

h4. Follow-ups (separate JIRAs)

* Virtual-thread processing mode (few pollers, one virtual thread per record, 
acknowledgements queued back to the poller) together with lock renewal 
({{RENEW}}), which only makes sense with asynchronous commits.
* Batching facade (one exchange per poll with per-record acknowledgement).
* camel-quarkus {{kafka-share}} extension; the camel-spring-boot starter is 
generated automatically.

_Claude Code on behalf of Croway_


> 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