Skip to content

NIFI-15961 - Add Kafka share-group support to ConsumeKafka - #11271

Open
pvillard31 wants to merge 1 commit into
apache:mainfrom
pvillard31:NIFI-15961
Open

NIFI-15961 - Add Kafka share-group support to ConsumeKafka#11271
pvillard31 wants to merge 1 commit into
apache:mainfrom
pvillard31:NIFI-15961

Conversation

@pvillard31

Copy link
Copy Markdown
Contributor

Summary

NIFI-15961 - Add Kafka share-group support to ConsumeKafka

Introduce share-group consumption for ConsumeKafka backed by Kafka 4.2's GA share consumer (KIP-932). A new "Group Type" property selects between the existing classic Consumer Group code path (default; preserves all prior behavior) and the new Share Group path, and a new "Acknowledgement Mode" property gates per-record vs. implicit acknowledgement when the share-group path is selected.

  • New KafkaShareConsumerService API contract plus Kafka4ShareConsumerService implementation wrapping org.apache.kafka.clients.consumer.ShareConsumer.
  • KafkaConnectionService gains a default getShareConsumerService() so third-party connection services remain ABI-compatible and only opt in when they support share groups.
  • Kafka3ConnectionService strips the classic-consumer properties rejected by ShareConsumerConfig.SHARE_GROUP_UNSUPPORTED_CONFIGS before constructing the share consumer.
  • ConsumeKafka.verify() exercises a share-group subscription and releases any sampled records back to the share group.
  • Bump the integration-test broker to apache/kafka:4.2.0
  • Add unit tests for the share-consumer service (poll, explicit/implicit ack, commit, rollback, close, tombstone handling), the processor's share-group selection/verification/ack-mode wiring, and an integration test ConsumeKafkaShareGroupIT.

Tracking

Please complete the following tracking steps prior to pull request creation.

Issue Tracking

Pull Request Tracking

  • Pull Request title starts with Apache NiFi Jira issue number, such as NIFI-00000
  • Pull Request commit message starts with Apache NiFi Jira issue number, as such NIFI-00000
  • Pull request contains commits signed with a registered key indicating Verified status

Pull Request Formatting

  • Pull Request based on current revision of the main branch
  • Pull Request refers to a feature branch with one commit containing changes

Verification

Please indicate the verification steps performed prior to pull request creation.

Build

  • Build completed using ./mvnw clean install -P contrib-check
    • JDK 21
    • JDK 25

Licensing

  • New dependencies are compatible with the Apache License 2.0 according to the License Policy
  • New dependencies are documented in applicable LICENSE and NOTICE files

Documentation

  • Documentation formatting appears as expected in rendered files
Introduce share-group consumption for ConsumeKafka backed by Kafka 4.2's
GA share consumer (KIP-932). A new "Group Type" property selects between
the existing classic Consumer Group code path (default; preserves all
prior behavior) and the new Share Group path, and a new "Acknowledgement
Mode" property gates per-record vs. implicit acknowledgement when the
share-group path is selected.

- New KafkaShareConsumerService API contract plus Kafka4ShareConsumerService
  implementation wrapping org.apache.kafka.clients.consumer.ShareConsumer.
- KafkaConnectionService gains a default getShareConsumerService() so
  third-party connection services remain ABI-compatible and only opt in
  when they support share groups.
- Kafka3ConnectionService strips the classic-consumer properties rejected
  by ShareConsumerConfig.SHARE_GROUP_UNSUPPORTED_CONFIGS before constructing
  the share consumer.
- ConsumeKafka.verify() exercises a share-group subscription and releases
  any sampled records back to the share group.
- Bump the integration-test broker to apache/kafka:4.2.0 (the version that
  matches the bundled kafka-clients and the first release where share
  groups are GA).
- Add unit tests for the share-consumer service (poll, explicit/implicit
  ack, commit, rollback, close, tombstone handling), the processor's
  share-group selection/verification/ack-mode wiring, and an integration
  test ConsumeKafkaShareGroupIT.

Signed-off-by: Pierre Villard <pierre.villard.fr@gmail.com>
@pvillard31

Copy link
Copy Markdown
Contributor Author

I rebased on latest and I'll add Kafka Share-Group backlog reporting as a follow-up item.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

1 participant