Versions Compared

Key

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

Table of Contents

Status

Current state: Under DiscussionAccepted

Discussion thread: here

JIRA:

Jira
serverASF JIRA
serverId5aa69414-a9e9-3523-82ec-879b028fb15b
keyKAFKA-19446

...

Code Block
{
  "apiKey": 27,
  "type": "request",
  "listeners": ["broker"],
  "name": "WriteTxnMarkersRequest",
  "validVersions": "1-2",
  "flexibleVersions": "1+",
  "fields": [
    { "name": "Markers", "type": "[]WritableTxnMarker", "versions": "0+",
      "about": "The transaction markers to be written.", "fields": [
      { "name": "ProducerId", "type": "int64", "versions": "0+",
        "entityType": "producerId", "about": "The current producer ID." },
      { "name": "ProducerEpoch", "type": "int16", "versions": "0+",
        "about": "The current epoch associated with the producer ID." },
      { "name": "TransactionResult", "type": "bool", "versions": "0+",
        "about": "The result (false = ABORT, true = COMMIT)." },
      { "name": "Topics", "type": "[]WritableTxnMarkerTopic", "versions": "0+",
        "about": "Each topic to write markers for.", "fields": [
        { "name": "Name", "type": "string", "versions": "0+",
          "entityType": "topicName", "about": "The topic name." },
        { "name": "PartitionIndexes", "type": "[]int32", "versions": "0+",
          "about": "Partition indexes to write markers for." }
      ]},
      { "name": "CoordinatorEpoch", "type": "int32", "versions": "0+",
        "about": "Epoch of the transaction state partition hosting this coordinator." },
      // --------- NEW FIELD (TV) ADDED BELOW --------
      { "name": "TransactionVersion", "type": "int8", "versions": "2+", "ignorable": true,
        "about": "Transaction version: 0/1 = legacy (TV0/TV1), 2 = TV2.", "default": "0" }
    ]}
  ]
}

...

Leaders must apply strict validation when TV is known to be >=TV2:

Code Block
// Pseudocode inside appendEndTxnMarker() or where checkProducerEpoch is called:

short current = updatedEntry.producerEpoch();
short marker = markerProducerEpoch;          
// Check if TransactionVersion field is available (version 2+)
int txnVersion = (request.version >= 2) ? request.transactionVersion() : 1;  

if (txnVersion >= 2) {
    // TV2: coordinator bumps epoch before marker; duplicates carry old epoch.
    // Accept only strictly greater epoch.
    if (marker <= current) {
        throw new InvalidProducerEpochException("Reject late/dup TV2 marker: " +
            "markerEpoch=" + marker + " <= currentEpoch=" + current);
    }
} else {
    // Legacy behavior
    if (marker < current) {
        throw new InvalidProducerEpochException("Marker epoch < current.");
    }
}

...