You are viewing an old version of this page. View the current version.

Compare with Current View Page History

« Previous Version 26 Next »

Status

Current state: Under discussion

Discussion thread:  here

JIRA: KAFKA-20367 - Getting issue details... STATUS

1 Motivation

1.1 The Problem: Scaling Kafka-to-Kafka Pipelines Today

Currently, Kafka Connect sink connectors rely on traditional consumer groups that enforce a strict 1:1 mapping between partitions and tasks. T

his model is often incompatible with unordered message processing and creates three primary bottlenecks for task queue workloads:

1. Partition-Coupled Scaling: Parallelism is hard-limited by the partition count

2. Head-of-Line Blocking: Because partition ownership is exclusive, a single slow task—often caused by downstream latency—stalls all subsequent records in its assigned partitions

3. Rebalance-Driven Gaps: Adding or removing tasks triggers "rebalance storms."

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: Decouples parallelism from partition count,
- No Head-of-Line Blocking: Supports unordered message processing; if a task slows down, records time out and are redelivered to available workers.
- Seamless Scaling: Eliminates "rebalance storms" by removing the partition assignment protocol, ensuring zero downtime during task membership changes.

Note: The share groups are only suitable for connectors with idempotent, order-independent processing.

2. Scope

2.1 In Scope (What we are building)

  • New Task Type: Introducing WorkerShareSinkTask to handle Share Group logic without changing existing connector code.

  • Flexible Activation: Toggle queue semantics globally or per-connector via consumer.override.group.protocol=share.

  • Delivery Guarantees: * At-least-once: Standard support for all sink types.

    • Exactly-once: Supported for same-cluster Kafka-to-Kafka paths (requires KIP-1289).

  • Observability: New Share Group metrics (acquisition, release, and rejection rates) integrated into the existing sink-task-metrics group.

2.2 Out of Scope (Future/Separate efforts)

  • Source Connectors: Share Group support is currently for Sinks only (Source support and MirrorMaker 2 are excluded).

  • Cross-Cluster EOS: Exactly-once delivery between different Kafka clusters is not supported in this phase.

  • API Changes: No modifications will be made to the public SinkTask Java API or individual connector codebases.

  • Complex Transactions: External 2PC coordinators and cross-cluster transactional protocols are not addressed.


3. Public Interfaces

3.1 New Configuration Properties

Worker-Level (connect-distributed.properties)

  • consumer.group.protocol: Set to share to enable KafkaShareConsumer globally for all sink tasks (Default: consumer).

Connector-Level (Per-connector JSON)

PropertyDefaultDescription
consumer.override.group.protocolInheritedSet to share to opt a specific connector into queue semantics.
share.group.idconnect-<name>Custom Share Group ID; follows standard naming conventions.
share.acknowledgement.modeexplicitexplicit: Acknowledge after task.put(). implicit: Acknowledge on the next poll.
share.acquisition.lock.timeout.ms30000Max time a record stays acquired before re-delivery. Must exceed task.put() latency.
share.delivery.semanticsat-least-onceToggle between at-least-once and exactly-once (requires KIP-1289).
share.max.delivery.count5Max re-delivery attempts before sending to a Dead Letter Queue.

3.2 New / Modified Java Interfaces


3.2.1 `WorkerShareSinkTask` (new class)


A new internal class—parallel to WorkerSinkTask—that drives SinkTask using a KafkaShareConsumer.

  • No API Changes: The public SinkTask interface and put() contract remain identical; existing connectors require no code modifications.

```
// New class: parallel to WorkerSinkTask but backed by ShareConsumer
class WorkerShareSinkTask extends WorkerTask<ConsumerRecord<byte[], byte[]>, SinkRecord> {
    private final ShareConsumer<byte[], byte[]> shareConsumer;
    private final SinkTask task;
    // ...
}
```

Note: The existing `SinkTask` interface is not modified. Connectors do not need code changes. The `put(Collection<SinkRecord>)` contract remains the same.

The difference is entirely in the worker runtime:

AspectWorkerSinkTask (Traditional)WorkerShareSinkTask (Proposed)
ConsumerKafkaConsumerKafkaShareConsumer
TrackingConsumer Offsets + commitSync()Per-record acknowledge(ACCEPT)
RebalanceRebalance listener triggers open/closeNone. task.open() called once at startup.
FailuresPause consumer and retry batchacknowledge(RELEASE) for broker re-delivery

3.2.2`Worker.baseConsumerConfigs()` (modified)

Updated to detect group.protocol=share. It dynamically constructs ShareConsumerConfig properties (like share.group.id) instead of traditional consumer configs.

Metrics

New sensors are registered only in Share Group mode to keep dashboards clean and verify the connector state. All metrics belong to the existing sink-task-metrics group.

New Share-Specific Metrics:

  • sink-record-acquire: Rate/Total of records pulled from the group.

  • sink-record-acknowledge: Rate/Total of successful ACCEPT acks.

  • sink-record-release/reject: Rate/Total of records released for retry or rejected to DLQ.

  • acknowledge-time: Time between poll() and acknowledge().

  • sink-record-redelivery: Total records with delivery count > 1.

Exclusions: The following traditional sensors are not registered in share mode as they are not applicable: partition-count, offset-seq-number, and offset-commit-completion.

Proposed Changes

At‑Least‑Once (Share Group → SinkTask → External Sink)

Exactly‑Once (Same‑Cluster Kafka‑to‑Kafka, KIP‑1289)



`WorkerShareSinkTask` Lifecycle

Initialization

```
void initialize() {
    // 1. Create KafkaShareConsumer with resolved configs
    this.shareConsumer = new KafkaShareConsumer<>(shareConsumerConfigs);
    
    // 2. Subscribe to configured topics
    List<String> topics = SinkConnectorConfig.parseTopicsList(taskConfig);
    shareConsumer.subscribe(topics);
    
    // 3. Open the task (no partition-level open/close with share groups)
    task.initialize(context);
    task.start(taskConfig);
}
```

Main Loop (iteration)


```
void iteration() {
    // 1. Poll records from share group
    ConsumerRecords<byte[], byte[]> records = shareConsumer.poll(Duration.ofMillis(pollTimeoutMs));
    
    if (records.isEmpty()) return;
    
    // 2. Convert to SinkRecords (same as today)
    List<SinkRecord> sinkRecords = convertMessages(records);
    
    // 3. Deliver to task
    try {
        task.put(sinkRecords);
        
        // 4a. Success: acknowledge all records as ACCEPT
        for (ConsumerRecord<byte[], byte[]> record : records) {
            shareConsumer.acknowledge(record, AcknowledgeType.ACCEPT);
        }
        
    } catch (RetriableException e) {
        // 4b. Retriable failure: RELEASE records for re-delivery
        for (ConsumerRecord<byte[], byte[]> record : records) {
            shareConsumer.acknowledge(record, AcknowledgeType.RELEASE);
        }
        log.warn("Retriable error, records released for re-delivery", e);
        
    } catch (Throwable t) {
        // 4c. Fatal failure: REJECT records (to DLQ if configured)
        for (ConsumerRecord<byte[], byte[]> record : records) {
            shareConsumer.acknowledge(record, AcknowledgeType.REJECT);
        }
        throw new ConnectException("Unrecoverable error", t);
    }
    
    // 5. Commit acknowledgments to broker
    if (shouldCommit()) {
        shareConsumer.commitSync();
    }
}
```

Ensuring No Data Loss (At-Least-Once)

The at-least-once guarantee is achieved through the following invariant:

> A record is acknowledged (ACCEPT) only after `task.put()` returns successfully.

If the task or worker crashes between `poll()` and `acknowledge()`:
- The record remains in ACQUIRED state on the broker
- The acquisition lock timer expires after `share.acquisition.lock.timeout.ms`
- The broker transitions the record back to AVAILABLE
- Another task acquires and processes it

If the worker crashes after `acknowledge(ACCEPT)` but before `commitSync()`:
- The implicit acknowledgment mode sends acks on the next `poll()`, so uncommitted acks may be lost
- The explicit mode (default) uses `commitSync()` which is durable. If the commit fails, the record stays in ACQUIRED and will time out and re-deliver.

Duplicate delivery can occur when a task successfully calls `task.put()` and `acknowledge(ACCEPT)` but crashes before the downstream system confirms persistence.

This is inherent to at-least-once semantics. Sink connectors targeting idempotent systems (databases with upsert, object stores with overwrite) naturally handle this.

Exactly-Once Semantics (Future Phase, requires KIP-1289)

For Kafka-to-Kafka pipelines (e.g., MirrorMaker2), exactly-once can be achieved by binding the Share Group acknowledgment to the producer's transaction:

```
// Exactly-once CTP pattern in WorkerShareSinkTask
void iterationExactlyOnce() {
    ConsumerRecords<byte[], byte[]> records = shareConsumer.poll(Duration.ofMillis(pollTimeoutMs));
    if (records.isEmpty()) return;
    
    producer.beginTransaction();
    
    try {
        // Produce transformed records to output topics
        for (SinkRecord record : convertMessages(records)) {
            producer.send(new ProducerRecord<>(outputTopic, record.key(), record.value()));
        }
        
        // Bind share acks to this transaction (KIP-1289)
        producer.sendShareAcksToTransaction(
            ShareAcknowledgements.fromRecords(records, AcknowledgeType.ACCEPT),
            shareConsumer.groupMetadata()
        );
        
        producer.commitTransaction();
        // Output records AND source acknowledgments commit atomically
        
    } catch (Exception e) {
        producer.abortTransaction();
        // Both output records AND source acknowledgments are rolled back
        // Records will be re-delivered by the broker
    }
}
```

 Configuration Resolution Order

```
Worker config (connect-distributed.properties)
    -> consumer.group.protocol=share          (global default)
    
Connector config (per-connector JSON)
    -> consumer.override.group.protocol=share  (per-connector override)
    -> share.group.id=my-custom-group          (explicit share group name)
    -> share.acknowledgement.mode=explicit     (ack behavior)
```

The existing `consumer.override.*` mechanism in Kafka Connect (governed by `connector.client.config.override.policy`) is reused. No new override mechanism is introduced.

Note: Share groups use a different state topic (__share_group_state), but looks like __consumer_offsets will be used for memebership, so if we do not delete the group before switching it can cause problem.

So user should make sure share group id is not equal to consumer group id at anytime. We can have a check/validation while implementing it. 

4. Compatibility, Deprecation, and Migration Plan

4.1 Impact on Existing Users

No impact by default. The default `group.protocol` remains `consumer` (standard consumer group). Existing connectors continue to work identically.
Opt-in only. Share Groups are enabled per-connector or per-worker via configuration.
No connector code changes required. The `SinkTask` interface is unchanged. Any existing sink connector works with Share Groups without modification.

4.2 Migration Path

1. Pre-requisite: Kafka broker version must support Share Groups (4.0+).
2. Enable at worker level: Set `consumer.group.protocol=share` in `connect-distributed.properties` to make all sink connectors use Share Groups.
3. Or enable per-connector: Set `consumer.override.group.protocol=share` in the connector config JSON.
4. Tune acquisition lock timeout: Set `share.acquisition.lock.timeout.ms` to a value greater than the expected `task.put()` latency. The default of 30 seconds is suitable for most workloads.
5. Monitor: Use the new `share-sink-task.*` metrics to observe acknowledgment patterns and re-delivery rates.

4.3 Rollback

To revert, remove the `group.protocol=share` configuration. The connector will resume using standard consumer groups.

Note that Share Groups and consumer groups maintain separate offset tracking, so the consumer group will resume from its last committed offset

(which may be behind the Share Group's position).

4.4 Deprecation

No existing features are deprecated. This is purely additive.

5. Test Plan

5.1 Unit Tests

1. `WorkerShareSinkTaskTest`: Tests the core poll-put-acknowledge loop using a `MockShareConsumer`.
   - Verify ACCEPT after successful `task.put()`
   - Verify RELEASE after `RetriableException`
   - Verify REJECT after unrecoverable exception
   - Verify `commitSync()` is called at configured intervals

2. `WorkerTest` (modified): Verify that `baseConsumerConfigs()` returns correct configs for `group.protocol=share`.

3. `SinkConnectorConfigTest` (modified): Validate the new configuration properties and their defaults.

5.2 Integration Tests

1. Basic Share Group Sink: Deploy a sink connector with `group.protocol=share` and verify all records are delivered.
2. Elastic Scaling: Start with 2 tasks, scale to 6, verify no records are lost and throughput increases.
3. Task Failure and Re-delivery: Kill a task mid-processing, verify records are re-delivered to surviving tasks within `acquisition.lock.timeout.ms`.
4. No Duplicate Loss: Produce N records, consume with at-least-once Share Group sink, verify received count >= N.
5. Interoperability: Verify that standard consumer group connectors and Share Group connectors can coexist in the same Connect cluster.

5.3 System Tests

1. Long-running throughput test: Measure throughput and latency of Share Group vs. consumer group sink connectors under sustained load.
2. Chaos test: Randomly kill tasks and brokers, verify zero data loss with at-least-once semantics.

6. Rejected Alternatives

Alternative 1: Modify the SinkTask Interface to Add acknowledge()

We considered adding `acknowledge(SinkRecord)` and `release(SinkRecord)` methods to the `SinkTask` interface, giving connectors explicit control over acknowledgments. This was rejected because:
- It would break backward compatibility with all existing sink connectors
- Most connectors don't need per-record acknowledgment control
- The worker runtime can make correct acknowledgment decisions based on `put()` success/failure

Alternative 2: Use Share Groups Only for MirrorMaker2

We considered limiting Share Group support to the `MirrorSourceConnector` only, as Kafka-to-Kafka is the most obvious use case. This was rejected because:
- It would require changes to the MM2 `consumer.assign()` model, which is complex
- Generic sink connectors (e.g., JDBC, Elasticsearch, S3) benefit equally from elastic scaling
- Building it into the Connect runtime benefits all connectors automatically

Alternative 3: Exactly-Once from Day One

We considered requiring exactly-once semantics for the initial implementation. This was rejected because:
- KIP-1289 (transactional share acknowledgments) is not yet implemented
- At-least-once is sufficient for the majority of sink connector use cases
- Idempotent sinks (upsert to database, overwrite to S3) achieve effective exactly-once with at-least-once delivery
- Exactly-once can be added as a follow-up without breaking changes

  • No labels