Versions Compared

Key

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

...

```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);
}
```

New Coordinator API

Reuses the existing TransactionCoordinator

new server-side component is a completeTransaction() method on ShareCoordinator, mirroring GroupCoordinator.completeTransaction().```java
public interface ShareTransactionCoordinator {
    void beginTransaction(String transactionId, Duration timeout);
    void prepareTransaction(String transactionId) throws TransactionPrepareException;
    void commitTransaction(String transactionId) throws TransactionCommitException;
    void abortTransaction(String transactionId);
    List<PendingTransaction> listPendingTransactions(String groupId);
}
```

New Configuration

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

...