From f85729c19bfd0940ce06a1da5f0f23653f9f92fd Mon Sep 17 00:00:00 2001 From: yuluo-yx Date: Sat, 8 Aug 2026 14:34:30 +0800 Subject: [PATCH] [ISSUE #10837] fix(client): skip truncated trace records --- .../rocketmq/client/trace/TraceDataEncoder.java | 10 +++++----- .../rocketmq/client/trace/TraceDataEncoderTest.java | 11 +++++++++++ 2 files changed, 16 insertions(+), 5 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/trace/TraceDataEncoder.java b/client/src/main/java/org/apache/rocketmq/client/trace/TraceDataEncoder.java index 1e66aa0498d..9530be42752 100644 --- a/client/src/main/java/org/apache/rocketmq/client/trace/TraceDataEncoder.java +++ b/client/src/main/java/org/apache/rocketmq/client/trace/TraceDataEncoder.java @@ -44,7 +44,7 @@ public static List decoderFromTraceDataString(String traceData) { String[] contextList = traceData.split(String.valueOf(TraceConstants.FIELD_SPLITOR)); for (String context : contextList) { String[] line = context.split(String.valueOf(TraceConstants.CONTENT_SPLITOR)); - if (line[0].equals(TraceType.Pub.name())) { + if (line.length >= 13 && line[0].equals(TraceType.Pub.name())) { TraceContext pubContext = new TraceContext(); pubContext.setTraceType(TraceType.Pub); pubContext.setTimeStamp(Long.parseLong(line[1])); @@ -77,7 +77,7 @@ public static List decoderFromTraceDataString(String traceData) { pubContext.setTraceBeans(new ArrayList<>(1)); pubContext.getTraceBeans().add(bean); resList.add(pubContext); - } else if (line[0].equals(TraceType.SubBefore.name())) { + } else if (line.length >= 8 && line[0].equals(TraceType.SubBefore.name())) { TraceContext subBeforeContext = new TraceContext(); subBeforeContext.setTraceType(TraceType.SubBefore); subBeforeContext.setTimeStamp(Long.parseLong(line[1])); @@ -91,7 +91,7 @@ public static List decoderFromTraceDataString(String traceData) { subBeforeContext.setTraceBeans(new ArrayList<>(1)); subBeforeContext.getTraceBeans().add(bean); resList.add(subBeforeContext); - } else if (line[0].equals(TraceType.SubAfter.name())) { + } else if (line.length >= 6 && line[0].equals(TraceType.SubAfter.name())) { TraceContext subAfterContext = new TraceContext(); subAfterContext.setTraceType(TraceType.SubAfter); subAfterContext.setRequestId(line[1]); @@ -112,7 +112,7 @@ public static List decoderFromTraceDataString(String traceData) { subAfterContext.setGroupName(line[8]); } resList.add(subAfterContext); - } else if (line[0].equals(TraceType.EndTransaction.name())) { + } else if (line.length >= 13 && line[0].equals(TraceType.EndTransaction.name())) { TraceContext endTransactionContext = new TraceContext(); endTransactionContext.setTraceType(TraceType.EndTransaction); endTransactionContext.setTimeStamp(Long.parseLong(line[1])); @@ -132,7 +132,7 @@ public static List decoderFromTraceDataString(String traceData) { endTransactionContext.setTraceBeans(new ArrayList<>(1)); endTransactionContext.getTraceBeans().add(bean); resList.add(endTransactionContext); - } else if (line[0].equals(TraceType.Recall.name())) { + } else if (line.length >= 7 && line[0].equals(TraceType.Recall.name())) { TraceContext recallContext = new TraceContext(); recallContext.setTraceType(TraceType.Recall); recallContext.setTimeStamp(Long.parseLong(line[1])); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/TraceDataEncoderTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/TraceDataEncoderTest.java index 26b7bda596b..9c7f1578edf 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/TraceDataEncoderTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/TraceDataEncoderTest.java @@ -64,6 +64,17 @@ public void testDecoderFromTraceDataString() { Assert.assertEquals(contexts.get(0).getTraceType(), TraceType.Pub); } + @Test + public void testDecoderSkipsTruncatedRecord() { + String truncatedRecord = TraceType.Pub.name() + TraceConstants.CONTENT_SPLITOR + + TraceConstants.FIELD_SPLITOR; + + List contexts = TraceDataEncoder.decoderFromTraceDataString(truncatedRecord + traceData); + + assertThat(contexts).hasSize(1); + assertThat(contexts.get(0).getTraceType()).isEqualTo(TraceType.Pub); + } + @Test public void testEncoderFromContextBean() { TraceContext context = new TraceContext();