This is an automated email from the ASF dual-hosted git repository.
chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 649481f2f9f KAFKA-20712 Fix transaction state name in log decoder
output (#22615)
649481f2f9f is described below
commit 649481f2f9f4f48c65f331ced48fcd7658ac6f8c
Author: Federico Valeri <[email protected]>
AuthorDate: Wed Jul 1 18:08:15 2026 +0200
KAFKA-20712 Fix transaction state name in log decoder output (#22615)
The TransactionLogMessageParser and TransactionLogMessageFormatter
output the numeric transactionStatus (e.g. 4) instead of the human
readable state name (e.g. CompleteCommit). This patch eesolve the byte
to the state name via TransactionState.fromId() in both the dump log
and consumer formatter code paths.
Reviewers: Luke Chen <[email protected]>, Justine Olshan
<[email protected]>, Chia-Ping Tsai <[email protected]>
---
.../java/org/apache/kafka/tools/DumpLogSegments.java | 15 ++++++++++++++-
.../tools/consumer/TransactionLogMessageFormatter.java | 16 +++++++++++++++-
.../java/org/apache/kafka/tools/DumpLogSegmentsTest.java | 4 ++--
.../apache/kafka/tools/consumer/ConsoleConsumerTest.java | 15 ++++++---------
.../consumer/TransactionLogMessageFormatterTest.java | 4 ++--
5 files changed, 39 insertions(+), 15 deletions(-)
diff --git a/tools/src/main/java/org/apache/kafka/tools/DumpLogSegments.java
b/tools/src/main/java/org/apache/kafka/tools/DumpLogSegments.java
index 8e6f3d8d096..f968aeea167 100644
--- a/tools/src/main/java/org/apache/kafka/tools/DumpLogSegments.java
+++ b/tools/src/main/java/org/apache/kafka/tools/DumpLogSegments.java
@@ -53,6 +53,8 @@ import
org.apache.kafka.coordinator.common.runtime.Deserializer;
import org.apache.kafka.coordinator.group.GroupCoordinatorRecordSerde;
import org.apache.kafka.coordinator.share.ShareCoordinatorRecordSerde;
import
org.apache.kafka.coordinator.transaction.TransactionCoordinatorRecordSerde;
+import org.apache.kafka.coordinator.transaction.TransactionState;
+import org.apache.kafka.coordinator.transaction.generated.TransactionLogValue;
import org.apache.kafka.metadata.MetadataRecordSerde;
import org.apache.kafka.metadata.bootstrap.BootstrapMetadata;
import org.apache.kafka.server.common.ApiMessageAndVersion;
@@ -725,8 +727,19 @@ public class DumpLogSegments {
@Override
protected JsonNode valueAsJson(ApiMessage message, short version) {
- return
org.apache.kafka.coordinator.transaction.generated.CoordinatorRecordJsonConverters
+ JsonNode json =
org.apache.kafka.coordinator.transaction.generated.CoordinatorRecordJsonConverters
.writeRecordValueAsJson(message, version);
+ if (message instanceof TransactionLogValue) {
+ byte statusId = ((TransactionLogValue)
message).transactionStatus();
+ String statusName;
+ try {
+ statusName = TransactionState.fromId(statusId).stateName();
+ } catch (IllegalStateException e) {
+ statusName = String.valueOf(statusId);
+ }
+ ((ObjectNode) json).put("transactionStatus", statusName);
+ }
+ return json;
}
}
diff --git
a/tools/src/main/java/org/apache/kafka/tools/consumer/TransactionLogMessageFormatter.java
b/tools/src/main/java/org/apache/kafka/tools/consumer/TransactionLogMessageFormatter.java
index cf9c540da3e..06d674f59b3 100644
---
a/tools/src/main/java/org/apache/kafka/tools/consumer/TransactionLogMessageFormatter.java
+++
b/tools/src/main/java/org/apache/kafka/tools/consumer/TransactionLogMessageFormatter.java
@@ -18,9 +18,12 @@ package org.apache.kafka.tools.consumer;
import org.apache.kafka.common.protocol.ApiMessage;
import
org.apache.kafka.coordinator.transaction.TransactionCoordinatorRecordSerde;
+import org.apache.kafka.coordinator.transaction.TransactionState;
import
org.apache.kafka.coordinator.transaction.generated.CoordinatorRecordJsonConverters;
+import org.apache.kafka.coordinator.transaction.generated.TransactionLogValue;
import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.node.ObjectNode;
public class TransactionLogMessageFormatter extends
CoordinatorRecordMessageFormatter {
public TransactionLogMessageFormatter() {
@@ -39,6 +42,17 @@ public class TransactionLogMessageFormatter extends
CoordinatorRecordMessageForm
@Override
protected JsonNode valueAsJson(ApiMessage message, short version) {
- return CoordinatorRecordJsonConverters.writeRecordValueAsJson(message,
version);
+ JsonNode json =
CoordinatorRecordJsonConverters.writeRecordValueAsJson(message, version);
+ if (message instanceof TransactionLogValue) {
+ byte statusId = ((TransactionLogValue)
message).transactionStatus();
+ String statusName;
+ try {
+ statusName = TransactionState.fromId(statusId).stateName();
+ } catch (IllegalStateException e) {
+ statusName = String.valueOf(statusId);
+ }
+ ((ObjectNode) json).put("transactionStatus", statusName);
+ }
+ return json;
}
}
diff --git
a/tools/src/test/java/org/apache/kafka/tools/DumpLogSegmentsTest.java
b/tools/src/test/java/org/apache/kafka/tools/DumpLogSegmentsTest.java
index 77fba98097c..629bf9b62d6 100644
--- a/tools/src/test/java/org/apache/kafka/tools/DumpLogSegmentsTest.java
+++ b/tools/src/test/java/org/apache/kafka/tools/DumpLogSegmentsTest.java
@@ -990,7 +990,7 @@ public class DumpLogSegmentsTest {
)),
Optional.of("{\"type\":\"0\",\"data\":{\"transactionalId\":\"txnId\"}}"),
Optional.of("{\"version\":\"0\",\"data\":{\"producerId\":123,\"producerEpoch\":0,\"transactionTimeoutMs\":0,"
+
-
"\"transactionStatus\":0,\"transactionPartitions\":[],\"transactionLastUpdateTimestampMs\":0,"
+
+
"\"transactionStatus\":\"Empty\",\"transactionPartitions\":[],\"transactionLastUpdateTimestampMs\":0,"
+
"\"transactionStartTimestampMs\":0}}")
);
@@ -1047,7 +1047,7 @@ public class DumpLogSegmentsTest {
)),
Optional.of("{\"type\":\"0\",\"data\":{\"transactionalId\":\"txnId\"}}"),
Optional.of("{\"version\":\"1\",\"data\":{\"producerId\":12,\"previousProducerId\":11,\"nextProducerId\":10,"
+
-
"\"producerEpoch\":2,\"transactionTimeoutMs\":14,\"transactionStatus\":0," +
+
"\"producerEpoch\":2,\"transactionTimeoutMs\":14,\"transactionStatus\":\"Empty\","
+
"\"transactionPartitions\":[{\"topic\":\"topic1\",\"partitionIds\":[0,1,2]}," +
"{\"topic\":\"topic2\",\"partitionIds\":[3,4,5]}],\"transactionLastUpdateTimestampMs\":123,"
+
"\"transactionStartTimestampMs\":13}}")
diff --git
a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java
b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java
index fcd3f5cacdd..642d6ed11a2 100644
---
a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java
+++
b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java
@@ -45,10 +45,9 @@ import
org.apache.kafka.coordinator.group.generated.OffsetCommitKey;
import
org.apache.kafka.coordinator.group.generated.OffsetCommitKeyJsonConverter;
import org.apache.kafka.coordinator.group.generated.OffsetCommitValue;
import
org.apache.kafka.coordinator.group.generated.OffsetCommitValueJsonConverter;
+import org.apache.kafka.coordinator.transaction.TransactionState;
import org.apache.kafka.coordinator.transaction.generated.TransactionLogKey;
import
org.apache.kafka.coordinator.transaction.generated.TransactionLogKeyJsonConverter;
-import org.apache.kafka.coordinator.transaction.generated.TransactionLogValue;
-import
org.apache.kafka.coordinator.transaction.generated.TransactionLogValueJsonConverter;
import org.apache.kafka.server.util.MockTime;
import com.fasterxml.jackson.databind.JsonNode;
@@ -297,7 +296,7 @@ public class ConsoleConsumerTest {
admin.createTopics(Set.of(newTopic));
produceMessagesWithTxn(cluster);
- String[] transactionLogMessageFormatter =
createConsoleConsumerArgs(cluster,
+ String[] transactionLogMessageFormatter =
createConsoleConsumerArgs(cluster,
Topic.TRANSACTION_STATE_TOPIC_NAME,
"org.apache.kafka.tools.consumer.TransactionLogMessageFormatter");
@@ -316,12 +315,10 @@ public class ConsoleConsumerTest {
assertNotNull(logKey);
assertEquals(transactionId, logKey.transactionalId());
- JsonNode valueNode = jsonNode.get("value");
- TransactionLogValue logValue =
-
TransactionLogValueJsonConverter.read(valueNode.get("data"),
TransactionLogValue.HIGHEST_SUPPORTED_VERSION);
- assertNotNull(logValue);
- assertEquals(0, logValue.producerId());
- assertEquals(0, logValue.transactionStatus());
+ JsonNode valueData = jsonNode.get("value").get("data");
+ assertNotNull(valueData);
+ assertEquals(0, valueData.get("producerId").asInt());
+ assertEquals(TransactionState.EMPTY.stateName(),
valueData.get("transactionStatus").asText());
} finally {
consumerWrapper.cleanup();
}
diff --git
a/tools/src/test/java/org/apache/kafka/tools/consumer/TransactionLogMessageFormatterTest.java
b/tools/src/test/java/org/apache/kafka/tools/consumer/TransactionLogMessageFormatterTest.java
index 7a6f86ace80..79aef4ce13c 100644
---
a/tools/src/test/java/org/apache/kafka/tools/consumer/TransactionLogMessageFormatterTest.java
+++
b/tools/src/test/java/org/apache/kafka/tools/consumer/TransactionLogMessageFormatterTest.java
@@ -60,7 +60,7 @@ public class TransactionLogMessageFormatterTest extends
CoordinatorRecordMessage
"data":{"producerId":100,
"producerEpoch":50,
"transactionTimeoutMs":500,
- "transactionStatus":4,
+ "transactionStatus":"CompleteCommit",
"transactionPartitions":[],
"transactionLastUpdateTimestampMs":1000,
"transactionStartTimestampMs":750}}}
@@ -75,7 +75,7 @@ public class TransactionLogMessageFormatterTest extends
CoordinatorRecordMessage
"data":{"producerId":100,
"producerEpoch":50,
"transactionTimeoutMs":500,
- "transactionStatus":4,
+ "transactionStatus":"CompleteCommit",
"transactionPartitions":[],
"transactionLastUpdateTimestampMs":1000,
"transactionStartTimestampMs":750}}}