DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
- 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
...