Versions Compared

Key

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

...

1. Worker acknowledges records → Records removed from Kafka
2. Checkpoint fails before sink write
3. Records lost (acknowledged but never persisted)

**Goal:** Enable exactly-once read semantics via transactional acknowledgements.###

Use Cases

- **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

...


New Consumer API Methods

```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

```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


New Wire Protocol Requests

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



Proposed Changes

Describe the new thing you want to do in appropriate detail. This may be fairly extensive and have large subsections of its own. Or it may be a few sentences. Use judgement based on the scope of the change.

1. Record State Machine

Introduce a new `LOCKED` state for records within a transaction:


Record State Machine:

AVAILABLE --[poll()]--> ACQUIRED
ACQUIRED --[txn.acknowledge()]--> LOCKED
ACQUIRED --[acquisitionLockTimeout, no txn]--> AVAILABLE
LOCKED --[commit()]--> ACKNOWLEDGED (removed)
LOCKED --[abort() or timeout]--> AVAILABLE (retry)


Key: `LOCKED` state is governed by **transaction timeout**, not acquisition lock.

2. Two-Phase Commit Protocol

Transaction States Transitions


- EMPTY -> ONGOING -> PREPARE -> COMMITTED (success)

- EMPTY -> ONGOING -> PREPARE -> ABORTED (abort after prepare)

- EMPTY -> ONGOING -> ABORTED (early abort)


Sequence Diagram

Two-Phase Commit Protocol Flow:

Step 1: Coordinator sends BeginBatchTxn(id) to Broker
Step 2: Executor-1 calls poll() to Broker - gets records
Step 3: Executor-2 calls poll() to Broker - gets records
Step 4: Executor-1 sends ack(txnId, records) to Broker - records become LOCKED
Step 5: Executor-2 sends ack(txnId, records) to Broker - records become LOCKED
Step 6: Coordinator sends PrepareBatchTxn(id) to Broker - state becomes PREPARE
Step 7: Broker responds "prepared" to Coordinator
Step 8: Coordinator sends CommitBatchTxn(id) to Broker - LOCKED records become ACKNOWLEDGED
Step 9: Broker responds "committed" to Coordinator



Image Added



3. Transaction ID Format

```
{groupId}-{transactionalIdPrefix}-{checkpointId}

Example: my-share-group-flink-job-42
```

4. State Storage

Extend `__share_group_state` topic with transaction records:

```java
public class ShareTransactionState {
    String transactionId;
    String groupId;
    TransactionState state;  // ONGOING, PREPARE, COMMITTED, ABORTED
    long createTimeMs;
    long prepareTimeMs;
    Map<TopicPartition, List<AcknowledgementBatch>> pendingAcks;
}
```

5. Failure Handling

Timeout Behavior

- ONGOING: Abort → Release records

- PREPARE: Coordinator decides (commit/abort)

- COMMITTED/ABORTED: Cleanup after retention

Recovery Protocol

```java
// On coordinator startup
List<PendingTransaction> pending = coordinator.listPendingTransactions(groupId);
for (PendingTransaction txn : pending) {
    if (txn.state == PREPARE && shouldCommit(txn)) {
        coordinator.commitTransaction(txn.id);  // Complete Phase 2
    } else {
        coordinator.abortTransaction(txn.id);   // Rollback
    }
}
```

Idempotency

All operations are idempotent by transaction ID:
- `beginTransaction(id)` - returns existing if already started
- `prepareTransaction(id)` - no-op if already prepared
- `commitTransaction(id)` - no-op if already committed
- `abortTransaction(id)` - no-op if already aborted/committed


6. Wire Protocol Details

ShareAcknowledgeTransactionalRequest

```
ShareAcknowledgeTransactionalRequest => GroupId TransactionId [Topics]
  GroupId => STRING
  TransactionId => STRING
  Topics => TopicId [Partitions]
    Partitions => Partition [AcknowledgementBatches]
      AcknowledgementBatches => FirstOffset LastOffset AcknowledgeType
```

ShareBeginTransactionRequest

```
ShareBeginTransactionRequest => GroupId TransactionId TimeoutMs
  GroupId => STRING
  TransactionId => STRING
  TimeoutMs => INT64
```

SharePrepareTransactionRequest

```
SharePrepareTransactionRequest => GroupId TransactionId
  GroupId => STRING
  TransactionId => STRING
```

ShareCommitTransactionRequest / ShareAbortTransactionRequest

```
ShareCommitTransactionRequest => GroupId TransactionId
  GroupId => STRING
  TransactionId => STRING
```

Compatibility, Deprecation, and Migration Plan

...

Impact on Existing Users

- **No breaking changes** for existing Share Group users
- `transactional` mode is opt-in via `share.acknowledgement.mode` config
- Existing `implicit` and `explicit` modes continue to work unchanged

Migration Path

1. Phase 1: Add `transactional` acknowledgement mode (backward compatible)
2. Phase 2: Streaming frameworks (Flink, Spark) implement transactional sources
3. Phase 3: Documentation and best practices for exactly-once semantics

Deprecation

- No deprecation of existing modes planned
- `transactional` mode recommended for exactly-once use cases

...

Test Plan

Describe in few sentences how the KIP will be tested. We are mostly interested in system tests (since unit-tests are specific to implementation details). How will we know that the implementation works as expected? How will we know nothing broke?

...