DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
Status
Current state: Draft
Discussion thread:
JIRA:
KAFKA-20367
-
Getting issue details...
STATUS
1 Motivation
1.1 The Problem: Scaling Kafka-to-Kafka Pipelines Today
Kafka Connect sink connectors consume from Kafka topics using traditional consumer groups.
In this model, each partition is exclusively assigned to one task. This creates two problems for Kafka-to-Kafka (and Kafka-to-external) pipelines:
1. Scaling is coupled to partition count.
If a topic has 12 partitions, you can run at most 12 sink tasks. I
ncreasing parallelism beyond the partition count requires repartitioning the source topic -- an operationally expensive and disruptive change.
2. Slow tasks block partitions.
If one sink task is slow (e.g., network latency to a downstream system), the records on its assigned partitions back up.
Other idle tasks cannot help because partition ownership is exclusive. This creates head-of-line blocking at the partition level.
3. Rebalance storms cause processing gaps.
When tasks are added, removed, or crash, consumer group rebalances revoke and reassign partitions.
During a rebalance, no task processes records from revoked partitions. With cooperative sticky rebalancing this is mitigated but not eliminated.
1.2 How Share Groups Solve This
Share Groups (KIP-932) introduce queue semantics for Kafka consumers. Unlike consumer groups, Share Groups do not assign partitions exclusively.
Instead, records from a partition are acquired by any available consumer in the group. After processing, the consumer acknowledges the record (ACCEPT, RELEASE, ARCHIEVED, or REJECT).
This provides:
- Elastic scaling independent of partition count.
50 tasks can process a 12-partition topic because records are distributed at the record level, not the partition level.
- No head-of-line blocking.
If one task is slow, acquired records time out and are re-delivered to another task.
- No rebalance disruption.
Share Groups have no partition assignment protocol. Adding or removing tasks does not trigger reassignment of partitions.
1.3 Kafka Connect as the Natural Integration Point
Kafka Connect is the standard framework for building data pipelines into and out of Kafka. Integrating Share Groups into Connect's sink connector
runtime gives every existing sink connector access to queue semantics with a configuration change only -- no connector code modifications required.
1.4 Delivery Semantics
| Mode | Guarantee | Mechanism |
|---|---|---|
| At-least-once(default) | Every record is delivered at least once; duplicates possible on failure | ShareConsumer.acknowledge(ACCEPT) after successful task.put(); records released on failure for re-delivery |
| Exactly-once | Every record is delivered exactly once | KIP-1289 transactional acknowledgments: producer.sendShareAcksToTransaction() binds acknowledgments to the output transaction |
At-least-once is the initial target. Exactly-once requires KIP-1289 (Transactional Acknowledgments for Share Groups) to be implemented and is described as a future phase.
1.5 Error Handling / Delivery Semantics
This KIP integrates with KIP-1191. In Share Group mode, Connect uses `AcknowledgeType.REJECT` for fatal errors and relies on the broker-side DLQ configured by KIP-1191. Retriable failures use `AcknowledgeType.RELEASE`. If a DLQ is configured for the share group, the broker is the single source of DLQ records; Connect does not emit its own DLQ records in this mode.
EO Constraints
For exactly-once, records with pending transactional acknowledgments must not be re-delivered while the transaction is open; KIP-1289 must suppress or renew acquisition locks until commit/abort. Operationally, `share.acquisition.lock.timeout.ms` must exceed worst-case `task.put()` plus transaction commit latency, otherwise duplicates are possible even with EOS.