This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new 1ed234725 fix(message): keep key trace lookups inside the queried key
(#5107)
1ed234725 is described below
commit 1ed234725e30c5a038c0e40f78e26cc96f0ff98e
Author: 风起 <[email protected]>
AuthorDate: Thu Oct 1 18:17:18 2026 +0800
fix(message): keep key trace lookups inside the queried key (#5107)
RocketMQ batches every trace context for one source topic into a single
trace message and indexes each business key on that message. Key lookup
parsed the whole body, so order-A included order-B and order-A-suffix.
Keep a context only when its keys column contains the query as a whole
token, using the same space separator as TraceDataEncoder. Message-id
lookup is unchanged.
Fixes #5106
---
.../provider/apache/RocketMQMessageProvider.java | 63 ++++++++++++++++++---
.../apache/RocketMQMessageProviderTest.java | 66 ++++++++++++++++++++++
2 files changed, 120 insertions(+), 9 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
index 66e052932..5ba5abd8a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
@@ -22,6 +22,7 @@ import org.apache.rocketmq.client.consumer.PullResult;
import org.apache.rocketmq.client.consumer.PullStatus;
import org.apache.rocketmq.client.trace.TraceConstants;
import org.apache.rocketmq.common.MixAll;
+import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageId;
@@ -654,7 +655,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
adminExt.queryMessage(effectiveTraceTopic(traceTopic),
msgId, TRACE_QUERY_MAX, begin, end);
if (traceResult != null && traceResult.getMessageList() != null) {
for (MessageExt traceMessage : traceResult.getMessageList()) {
- parseTraceBody(traceMessage.getBody(), msgId, nodes,
consumerStatus, true);
+ parseTraceBody(traceMessage.getBody(), msgId, null, nodes,
consumerStatus, true);
}
}
} catch (BusinessException e) {
@@ -679,10 +680,12 @@ public class RocketMQMessageProvider implements
MessageProvider {
}
/**
- * Trace lookup by business key. The key query already scopes the returned
trace messages to
- * the requested message, so the body parser does not filter on a message
id. The original
- * message topic is not required to query the global trace topic but is
kept in the signature
- * for API symmetry and logged for diagnostics.
+ * Trace lookup by business key. RocketMQ appends every trace context for
one source topic
+ * into the same trace message and indexes each business key on that
message, so the query
+ * result can contain contexts for other keys. Contexts are kept only when
their keys column
+ * contains {@code key} as a whole token. The original message topic is
not required to query
+ * the global trace topic but is kept in the signature for API symmetry
and logged for
+ * diagnostics.
*/
private TraceRecordVO getMessageTraceByKey(String instanceId,
DefaultMQAdminExt adminExt, String key,
String topic, String
traceTopic) {
@@ -701,7 +704,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
adminExt.queryMessage(effectiveTraceTopic(traceTopic),
key, TRACE_QUERY_MAX, begin, end);
if (traceResult != null && traceResult.getMessageList() != null) {
for (MessageExt traceMessage : traceResult.getMessageList()) {
- parseTraceBody(traceMessage.getBody(), null, nodes,
consumerStatus, false);
+ parseTraceBody(traceMessage.getBody(), null, key, nodes,
consumerStatus, false);
}
}
} catch (BusinessException e) {
@@ -760,10 +763,11 @@ public class RocketMQMessageProvider implements
MessageProvider {
* Parse a trace message body. Trace contexts are separated by STX ({@code
\u0002}) and the
* fields in each context are separated by SOH ({@code \u0001}); the first
field is the trace
* type. When {@code filterByMsgId} is true only contexts whose message id
matches
- * {@code targetMsgId} are kept; otherwise every context is parsed (used
by key lookups where
- * the query already scoped the trace messages to the requested key).
+ * {@code targetMsgId} are kept. When {@code targetKey} is non-null, only
contexts whose keys
+ * column contains that key as a whole token are kept. Message-id lookups
pass a null key and
+ * are not filtered by key.
*/
- private void parseTraceBody(byte[] body, String targetMsgId,
List<TraceNodeVO> nodes,
+ private void parseTraceBody(byte[] body, String targetMsgId, String
targetKey, List<TraceNodeVO> nodes,
List<ConsumerStatusVO> consumerStatus, boolean
filterByMsgId) {
if (body == null || body.length == 0) {
return;
@@ -784,6 +788,9 @@ public class RocketMQMessageProvider implements
MessageProvider {
if (filterByMsgId && !targetMsgId.equals(field(fields,
msgIdIndex))) {
continue;
}
+ if (targetKey != null && !traceKeysContain(traceType, fields,
targetKey)) {
+ continue;
+ }
try {
switch (traceType) {
case "Pub":
@@ -809,6 +816,44 @@ public class RocketMQMessageProvider implements
MessageProvider {
}
}
+ /**
+ * Keys column of a RocketMQ 5.5.0 trace context. Recall has no keys
column, so a key lookup
+ * cannot attribute it and the context is omitted. Pub, EndTransaction,
and SubBefore store
+ * keys at index 7; SubAfter stores them at index 5.
+ */
+ private static int traceKeysIndex(String traceType) {
+ return switch (traceType) {
+ case "Pub", "EndTransaction", "SubBefore" -> 7;
+ case "SubAfter" -> 5;
+ default -> -1;
+ };
+ }
+
+ /**
+ * Same token split {@code TraceDataEncoder} uses when it indexes a trace
message:
+ * {@code keys.split(MessageConst.KEY_SEPARATOR)} (a single space). {@code
order-A} matches
+ * {@code extra order-A} and does not match {@code order-A-suffix}.
+ */
+ private static boolean traceKeysContain(String traceType, String[] fields,
String queryKey) {
+ if (!StringUtils.hasText(queryKey)) {
+ return false;
+ }
+ int keysIndex = traceKeysIndex(traceType);
+ if (keysIndex < 0) {
+ return false;
+ }
+ String keys = field(fields, keysIndex);
+ if (!StringUtils.hasText(keys)) {
+ return false;
+ }
+ for (String token : keys.split(MessageConst.KEY_SEPARATOR)) {
+ if (queryKey.equals(token)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
// Pub layout (RocketMQ 5.5.0 TraceDataEncoder):
// type, time, region, group, topic, msgId,
// tags, keys, storeHost, bodyLength, costTime, msgType,
offsetMsgId, isSuccess
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
index 87494e274..d2d470831 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
@@ -44,6 +44,7 @@ import
org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
+import org.apache.rocketmq.studio.instance.message.ConsumerStatusVO;
import org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
@@ -1024,6 +1025,71 @@ class RocketMQMessageProviderTest {
assertThat(keyCaptor.getValue()).isEqualTo("shared-key");
}
+ @Test
+ void getMessageTraceByKeyKeepsOnlyContextsWhoseKeysContainTheQueryToken()
throws Exception {
+ // One trace message batches every context for the source topic. The
index lists each
+ // space-separated business key, so a lookup for order-A hits contexts
for order-B and
+ // order-A-suffix too. Only whole key tokens may be returned.
+ String pubA = traceContext("Pub", "1000", "cn", "producer-a",
"orders", "msg-a",
+ "tag", "order-A", "127.0.0.1:10911", "10", "5", "0",
"offset-msg-a", "true");
+ String pubMulti = traceContext("Pub", "1001", "cn", "producer-a2",
"orders", "msg-a2",
+ "tag", "extra order-A", "127.0.0.1:10911", "10", "5", "0",
"offset-msg-a2", "true");
+ String pubB = traceContext("Pub", "1002", "cn", "producer-b",
"orders", "msg-b",
+ "tag", "order-B", "127.0.0.1:10911", "10", "5", "0",
"offset-msg-b", "true");
+ String subA = traceContext("SubAfter", "req-a", "msg-a", "5", "true",
"order-A",
+ "0", "3000", "consumer-a");
+ String subB = traceContext("SubAfter", "req-b", "msg-b", "5", "false",
"order-B",
+ "0", "3001", "consumer-b");
+ String subSuffix = traceContext("SubAfter", "req-suffix",
"msg-prefix", "5", "false",
+ "order-A-suffix", "0", "3002", "consumer-prefix");
+ String txA = traceContext("EndTransaction", "1003", "cn", "tx-a",
"orders", "msg-a",
+ "tag", "extra order-A", "127.0.0.1:10911", "0", "tx-1",
"COMMIT_MESSAGE", "false");
+ String txB = traceContext("EndTransaction", "1004", "cn", "tx-b",
"orders", "msg-b",
+ "tag", "order-B", "127.0.0.1:10911", "0", "tx-2",
"ROLLBACK_MESSAGE", "false");
+ String recall = traceContext("Recall", "2500", "cn",
"producer-recall", "orders", "msg-recall", "false");
+ MessageExt traceMessage = new MessageExt();
+ traceMessage.setBody(traceBody(pubA, pubMulti, pubB, subA, subB,
subSuffix, txA, txB, recall)
+ .getBytes(StandardCharsets.UTF_8));
+ when(adminExt.queryMessage(anyString(), anyString(), anyInt(),
anyLong(), anyLong()))
+ .thenReturn(new QueryResult(0L, List.of(traceMessage)));
+
+ TraceRecordVO byKey = provider.getMessageTraceByKey("instance-a",
"order-A", "orders", null);
+
+ assertThat(byKey.getNodes()).extracting(TraceNodeVO::getDescription)
+ .containsExactly(
+ "producer=producer-a, storeHost=127.0.0.1:10911",
+ "producer=producer-a2, storeHost=127.0.0.1:10911",
+ "group=consumer-a, contextCode=0",
+ "group=tx-a, transactionState=COMMIT_MESSAGE");
+
assertThat(byKey.getConsumerStatus()).extracting(ConsumerStatusVO::getGroup)
+ .containsExactly("consumer-a");
+
+ TraceRecordVO bySuffix = provider.getMessageTraceByKey(
+ "instance-a", "order-A-suffix", "orders", null);
+
+ assertThat(bySuffix.getNodes()).extracting(TraceNodeVO::getDescription)
+ .containsExactly("group=consumer-prefix, contextCode=0");
+
assertThat(bySuffix.getConsumerStatus()).extracting(ConsumerStatusVO::getGroup)
+ .containsExactly("consumer-prefix");
+
+ // Message-id filtering stays exact and is not replaced by the key
token check.
+ TraceRecordVO byMsgB = provider.getMessageTrace("instance-a", "msg-b",
"orders");
+
+ assertThat(byMsgB.getNodes()).extracting(TraceNodeVO::getDescription)
+ .containsExactly(
+ "producer=producer-b, storeHost=127.0.0.1:10911",
+ "group=consumer-b, contextCode=0",
+ "group=tx-b, transactionState=ROLLBACK_MESSAGE");
+
assertThat(byMsgB.getConsumerStatus()).extracting(ConsumerStatusVO::getGroup)
+ .containsExactly("consumer-b");
+
+ TraceRecordVO byRecall = provider.getMessageTrace("instance-a",
"msg-recall", "orders");
+
+
assertThat(byRecall.getNodes()).extracting(TraceNodeVO::getTitle).containsExactly("recall");
+
assertThat(byRecall.getNodes().get(0).getDescription()).contains("producer-recall");
+ assertThat(byRecall.getConsumerStatus()).isEmpty();
+ }
+
@Test
void getMessageTraceByKeyUsesDefaultTraceTopicWhenNotSpecified() throws
Exception {
when(adminExt.queryMessage(anyString(), anyString(), anyInt(),
anyLong(), anyLong()))