Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

...

- Apache Flink: Exactly-once checkpointing with Share Group sources
Apache Spark: Structured Streaming with Share Group consumers
Any coordinator-worker streaming framework requiring atomic acknowledgements

Public Interfaces

...

sendShareAcksToTransaction API

...

Method

KafkaProducer.sendShareAcksToTransaction(acks, groupId)

```java
public interface ShareConsumer<K, V> {
    // Existing
    void acknowledge(ConsumerRecord<K, V> record, AcknowledgeType type);
    
    // New: Transactional acknowledgement
    void acknowledge(ConsumerRecord<K, V> record, AcknowledgeType type, 
                     String transactionId);
}
```

...

Coordinator API

Reuses the existing TransactionCoordinator

...

Config                           Default      Description
share.acknowledgement.mode       explicit     Values: implicit, explicit, transactional
share.transaction.timeout.ms     60000        Maximum transaction duration before auto-abort
share.transaction.max.pending    5            Maximum concurrent transactions per group


New

...

Request                          Description
ShareAcknowledgeTransactional    Acknowledge records within a transaction
ShareBeginTransaction            Start a new batch transaction
SharePrepareTransaction          Phase 1: Prepare transaction
ShareCommitTransaction           Phase 2: Commit transaction
ShareAbortTransaction            Abort transaction and release records

New Metrics

Metric                                Type        Description
share-transaction-active              Gauge       Active transactions
share-transaction-prepare-time-ms     Histogram   Prepare latency
share-transaction-commit-time-ms      Histogram   Commit latency
share-transaction-abort-total         Counter     Aborted transactions
share-transaction-timeout-total       Counter     Timed-out transactions

...

Share acks are written transactionally to __share_group_state using CoordinatorRuntime.scheduleTransactionalWriteOperation().

...

3. Wire Protocol Details

TxnShareAcknowledgeRequest

...

mirrors AddOffsetsToTxnRequest; This tells the TransactionCoordinator to add __share_group_state partitions to the producer's transaction.

Consume-Transform-Produce Pattern


Compatibility, Deprecation, and Migration Plan

...