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

...

APIs 

  • KafkaProducer.sendShareAcksToTransaction(acks, groupId) - for CTP 
    • This mirrors the existing KafkaProducer.sendOffsetsToTransaction(offsets, groupMetadata)

```

// When you have a producer and want atomic output + acks

producer.beginTransaction();

producer.send(output);

producer.sendShareAcksToTransaction(acks, groupId);

producer.commitTransaction();

```


  • TransactionalShareAcknowledger - for standalone ack transactions

```

// When there's no producer (Flink source, Spark source, pure consumer)

public class TransactionalShareAcknowledger implements Closeable {

    public TransactionalShareAcknowledger(Map<String, Object> config);

    public void initTransactions();

    public void commitAcknowledgements(

        Map<TopicPartition, List<ShareAcknowledgement>> acks, String groupId);

    public void abortAcknowledgements();

    public void close();

}

```

Internally, TransactionalShareAcknowledger is a thin wrapper around a KafkaProducer (or its TransactionManager). It use the exact same RPCs - InitProducerId, AddShareAcksToTxn, TxnShareAcknowledge, EndTxn.

No new server-side infrastructure needed.

Coordinator API

  • Reuses the existing TransactionCoordinator
  • new server-side component is a completeTransaction() method on ShareCoordinator, mirroring GroupCoordinator.completeTransaction().
  • WriteTxnMarkers as the mechanism for the transaction coordinator to tell the group coordinator to complete transactional operations on __consumer_offsets
    • For consumer group we have => groupCoordinator.completeTransaction(partition, ...)
    • Similarly implement shareCoordinator.completeTransaction(partition, ...) for share group 

...