Croway opened a new pull request, #27495: URL: https://github.com/apache/camel/pull/27495
> [!IMPORTANT] > **Draft, stacked on #27488.** This branch contains the commit of #27488 (`camel-kafka-common`). Only the last commit (`CAMEL-25029: camel-kafka-share - consumer for Kafka share groups (KIP-932)`) belongs to this PR. Once #27488 is merged on `main`, this branch will be rebased onto `main` and marked ready for review. ## Description [CAMEL-25029](https://issues.apache.org/jira/browse/CAMEL-25029): third of three PRs for the Kafka share group consumer (KIP-932, Queues for Kafka). 1. #27485 (merged): extract a reusable client configuration layer inside `camel-kafka`. 2. #27488: move that layer into `camel-kafka-common`. 3. **This PR:** add the `kafka-share` component in a new `camel-kafka-share` module, which depends only on `camel-kafka-common`. ### The component A consumer-only scheme, `kafka-share:topic`, with `groupId` (the share group) required. Several consumers can read the same partition, and each record is acknowledged on its own; the broker redelivers released records and counts deliveries. Producing stays on `kafka:`. | Exchange outcome | Acknowledgement | |---|---| | The route set `CamelKafkaShareAcknowledge` (`ACCEPT`, `RELEASE`, `REJECT`) | the header value | | Completed, or failed with an exception the route handled | `ACCEPT` | | Failed, or rolled back | `onFailure`, default `RELEASE` | | Not processed because the consumer is stopping | `RELEASE` | - Each of the `consumersCount` threads owns one `KafkaShareConsumer`, which runs in explicit acknowledgement mode. Every record of a poll is acknowledged on the polling thread before the next poll. - After the records of a poll are processed, their acknowledgements are committed. With `commitMode=SYNC` (the default), per-partition commit errors go to the exception handler; with `ASYNC` they are logged. - The delivery count is in the `CamelKafkaShareDeliveryCount` header. The topic, partition, offset, key, timestamp and headers use the same header names as `kafka:` (`KafkaShareConstants`). - `CamelKafkaShareAcknowledge` is a `Camel*` header, so the inbound `KafkaHeaderFilterStrategy` removes it from Kafka records: only the route can set it. - The brokers, security, deserializer, header and additional-properties options are inherited from `KafkaClientConfiguration`. `KafkaShareConfiguration` adds the share options and never writes the client properties that `ShareConsumerConfig` rejects (`auto.offset.reset`, `enable.auto.commit`, `group.protocol`, ...). A unit test checks this by building a real `KafkaShareConsumer` from the properties. - `KafkaShareClientFactory` is a separate factory interface, so the existing `KafkaClientFactory` of `kafka:` and camel-quarkus's implementation of it are untouched. - Other features: - a readiness health check (ready when every share consumer is created and subscribed); - `pollOnError` (`ERROR_HANDLER`, `RECONNECT`, `STOP`, `DISCARD`/`RETRY`); - the create-consumer backoff of the component; - consumers are started only once the CamelContext is started. - `AbstractKafkaComponent.pendingConsumer` is now public, as the share component lives in another package. Docs: `kafka-share-component.adoc` covers the broker requirements, the acknowledgement table, the delivery-count example from the JIRA issue, what `consumersCount` means for a share group, at-least-once delivery and the differences from `kafka:`. A short "Share groups (queues)" pointer is added to `kafka-component.adoc`. **Not in this PR (follow-ups):** - a manual acknowledgement API for asynchronous hand-off; - `RENEW` of the acquisition lock; - a dev console; - a batching mode; - a virtual-thread processing mode; - a camel-quarkus extension and a kamelet. ### test-infra `KafkaContainer` configures the broker through environment variables, so the `server.properties` of the image (which sets `share.coordinator.state.topic.replication.factor=1`) is not used. The broker default is 3, so on the single test broker the share coordinator cannot create its state topic, and share consumers receive no records without any error. The three places that create the test container now set the replication factor and min ISR of that topic to 1. It only affects the internal topic of share groups, so the existing Kafka tests are not affected. ### Verification - Unit tests (11), using a `MockShareConsumer` that records acknowledgements: an acknowledgement for each outcome in the table, `onFailure=REJECT`, the mapping of records to messages and headers, the generated client properties (accepted by `KafkaShareConsumer`, and a rejected key in `additionalProperties` fails), and endpoint options. - `KafkaShareConsumerIT` against the test-infra broker (`apache/kafka:4.3.1`): - 60 records on one partition are shared by 3 consumers with no duplicates; - a failed record is redelivered with `CamelKafkaShareDeliveryCount=2`; - a record the route rejects is not redelivered. - `mvnd clean install -DskipTests` from the root: green, and the regenerated catalog, DSL, BOM, docs navigation and camel-util files are committed. ## Target - [x] I checked that the commit is targeting the correct branch (Camel 4 uses the `main` branch) ## Tracking - [x] If this is a large change, bug fix, or code improvement, I checked there is a [JIRA issue](https://issues.apache.org/jira/browse/CAMEL) filed for the change (usually before you start working on it). ## Apache Camel coding standards and style - [x] I checked that each commit in the pull request has a meaningful subject line and body. - [x] I have run `mvn clean install -DskipTests` locally from root folder and I have committed all auto-generated changes. _Claude Code on behalf of Croway_ -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
