From 78e993d5bb3ae95b9ce4cd4199de5e24cd19294c Mon Sep 17 00:00:00 2001 From: DaxPay Dev Date: Tue, 23 Jun 2026 10:09:55 +0800 Subject: [PATCH] =?UTF-8?q?feat(demo):=20=E6=B6=88=E6=81=AF=E9=98=9F?= =?UTF-8?q?=E5=88=97=E6=BC=94=E7=A4=BA=E5=A2=9E=E5=8A=A0=20traceId=20?= =?UTF-8?q?=E8=B4=AF=E7=A9=BF=E5=B1=95=E7=A4=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - DemoArtemisMessage 增加 producerTraceId 字段, 随消息体跨 MQ 传递 - DemoMessageResult 增加 producerTraceId/consumerTraceId/tracePropagated 字段 - Controller 发送时注入 MDC traceId 到消息体 - 3 个 Consumer 处理时读 MDC traceId 并与生产者对比 - 页面可直观看到 trace 是否贯穿(全程采样 1.0) --- .../artemis/consumer/DemoDelayConsumer.java | 9 ++++--- .../artemis/consumer/DemoQueueConsumer.java | 10 ++++++-- .../artemis/consumer/DemoTopicConsumer.java | 10 ++++++-- .../controller/ArtemisDemoController.java | 7 ++++-- .../artemis/model/DemoArtemisMessage.java | 3 +++ .../artemis/result/DemoMessageResult.java | 25 ++++++++++++++++--- 6 files changed, 52 insertions(+), 12 deletions(-) diff --git a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoDelayConsumer.java b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoDelayConsumer.java index 2769d52b6..be941bad4 100644 --- a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoDelayConsumer.java +++ b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoDelayConsumer.java @@ -6,6 +6,8 @@ import cn.daxpay.open.demo.artemis.constants.DemoArtemisConstants; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import cn.daxpay.open.platform.common.json.util.JacksonUtil; +import cn.daxpay.open.platform.core.code.CommonCode; +import org.slf4j.MDC; import org.springframework.jms.annotation.JmsListener; import org.springframework.stereotype.Component; @@ -32,9 +34,10 @@ public class DemoDelayConsumer { log.warn("Delay 消息解析失败,忽略: json={}, error={}", json, e.getMessage()); return; } - DemoMessageResult result = DemoMessageResult.from(message, "demo-delay-consumer"); + DemoMessageResult result = DemoMessageResult.from(message, "demo-delay-consumer", MDC.get(CommonCode.TRACE_ID)); store.add(result); - log.info("Delay 消费成功: id={}, content={}, costMillis={}", - message.getId(), message.getContent(), result.getCostMillis()); + log.info("Delay 消费成功: id={}, content={}, costMillis={}, producerTraceId={}, consumerTraceId={}", + message.getId(), message.getContent(), result.getCostMillis(), + message.getProducerTraceId(), result.getConsumerTraceId()); } } diff --git a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoQueueConsumer.java b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoQueueConsumer.java index 510c05ca7..e4e4d4fc8 100644 --- a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoQueueConsumer.java +++ b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoQueueConsumer.java @@ -6,6 +6,8 @@ import cn.daxpay.open.demo.artemis.constants.DemoArtemisConstants; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import cn.daxpay.open.platform.common.json.util.JacksonUtil; +import cn.daxpay.open.platform.core.code.CommonCode; +import org.slf4j.MDC; import org.springframework.jms.annotation.JmsListener; import org.springframework.stereotype.Component; @@ -30,8 +32,12 @@ public class DemoQueueConsumer { log.warn("Queue 消息解析失败,忽略: json={}, error={}", json, e.getMessage()); return; } - DemoMessageResult result = DemoMessageResult.from(message, "demo-queue-consumer"); + // OTel JMS observation 若生效, MDC 已被注入与 producer 相同的 traceId + String consumerTraceId = MDC.get(CommonCode.TRACE_ID); + DemoMessageResult result = DemoMessageResult.from(message, "demo-queue-consumer", consumerTraceId); store.add(result); - log.info("Queue 消费成功: id={}, content={}", message.getId(), message.getContent()); + log.info("Queue 消费成功: id={}, content={}, producerTraceId={}, consumerTraceId={}", + message.getId(), message.getContent(), + message.getProducerTraceId(), consumerTraceId); } } diff --git a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoTopicConsumer.java b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoTopicConsumer.java index 27da8c503..a89b79e3b 100644 --- a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoTopicConsumer.java +++ b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/consumer/DemoTopicConsumer.java @@ -7,6 +7,8 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import cn.daxpay.open.platform.common.artemis.ArtemisBeanNames; import cn.daxpay.open.platform.common.json.util.JacksonUtil; +import cn.daxpay.open.platform.core.code.CommonCode; +import org.slf4j.MDC; import org.springframework.jms.annotation.JmsListener; import org.springframework.stereotype.Component; @@ -45,8 +47,12 @@ public class DemoTopicConsumer { log.warn("Topic 消息解析失败,忽略: json={}, error={}", json, e.getMessage()); return; } - DemoMessageResult result = DemoMessageResult.from(message, consumer); + // OTel JMS observation 若生效, MDC 已被注入与 producer 相同的 traceId + String consumerTraceId = MDC.get(CommonCode.TRACE_ID); + DemoMessageResult result = DemoMessageResult.from(message, consumer, consumerTraceId); store.add(result); - log.info("Topic 消费成功 [{}]: id={}, content={}", consumer, message.getId(), message.getContent()); + log.info("Topic 消费成功 [{}]: id={}, content={}, producerTraceId={}, consumerTraceId={}", + consumer, message.getId(), message.getContent(), + message.getProducerTraceId(), consumerTraceId); } } diff --git a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/controller/ArtemisDemoController.java b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/controller/ArtemisDemoController.java index fb83ebe99..ac4ed67ac 100644 --- a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/controller/ArtemisDemoController.java +++ b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/controller/ArtemisDemoController.java @@ -11,9 +11,11 @@ import lombok.RequiredArgsConstructor; import cn.daxpay.open.platform.common.artemis.service.ArtemisTemplateService; import cn.daxpay.open.platform.common.json.util.JacksonUtil; import cn.daxpay.open.platform.core.annotation.IgnoreAuth; +import cn.daxpay.open.platform.core.code.CommonCode; import cn.daxpay.open.platform.core.exception.BizInfoException; import cn.daxpay.open.platform.core.rest.Res; import cn.daxpay.open.platform.core.rest.result.Result; +import org.slf4j.MDC; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PostMapping; @@ -49,12 +51,13 @@ public class ArtemisDemoController { // 校验场景相关必填字段 validateSceneParam(param, scene); - // 构造消息体(公共字段) + // 构造消息体(公共字段), 注入当前 HTTP 请求的 traceId 用于跨 MQ 追踪贯穿验证 DemoArtemisMessage message = new DemoArtemisMessage() .setId(UUID.randomUUID().toString(true)) .setContent(param.getContent()) .setScene(scene.name()) - .setSendTime(OffsetDateTime.now()); + .setSendTime(OffsetDateTime.now()) + .setProducerTraceId(MDC.get(CommonCode.TRACE_ID)); // 序列化为 JSON 字符串,发送层只负责搬运文本,不参与对象转换 String json = JacksonUtil.toJson(message, false); diff --git a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/model/DemoArtemisMessage.java b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/model/DemoArtemisMessage.java index cac8dede7..cf538ca10 100644 --- a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/model/DemoArtemisMessage.java +++ b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/model/DemoArtemisMessage.java @@ -31,4 +31,7 @@ public class DemoArtemisMessage { /// 发送时间(UTC) private OffsetDateTime sendTime; + + /// 生产者 traceId(随消息体跨 MQ 传递, 与消费端 MDC traceId 对比可验证 trace 是否贯穿) + private String producerTraceId; } diff --git a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/result/DemoMessageResult.java b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/result/DemoMessageResult.java index f781540ec..c7ff66a9a 100644 --- a/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/result/DemoMessageResult.java +++ b/daxpay-demo/src/main/java/cn/daxpay/open/demo/artemis/result/DemoMessageResult.java @@ -1,5 +1,6 @@ package cn.daxpay.open.demo.artemis.result; import cn.daxpay.open.demo.artemis.model.DemoArtemisMessage; +import cn.hutool.core.util.StrUtil; import io.swagger.v3.oas.annotations.media.Schema; import lombok.Data; import lombok.experimental.Accessors; @@ -43,16 +44,34 @@ public class DemoMessageResult { @Schema(description = "消费者名称") private String consumer; - /// 从消息体构造消费记录,自动填充消费时间与耗时 - public static DemoMessageResult from(DemoArtemisMessage message, String consumer) { + /// 生产者 traceId(发送消息时的 HTTP traceId, 来自消息体) + @Schema(description = "生产者traceId") + private String producerTraceId; + + /// 消费者 traceId(@JmsListener 处理时的 MDC traceId, 由 OTel 自动注入) + @Schema(description = "消费者traceId") + private String consumerTraceId; + + /// trace 是否贯穿 MQ(producerTraceId 与 consumerTraceId 非空且相等) + @Schema(description = "trace是否贯穿MQ") + private boolean tracePropagated; + + /// 从消息体构造消费记录, 自动填充消费时间/耗时/traceId 贯穿信息 + public static DemoMessageResult from(DemoArtemisMessage message, String consumer, String consumerTraceId) { OffsetDateTime now = OffsetDateTime.now(); + String producerTraceId = message.getProducerTraceId(); + boolean propagated = StrUtil.isNotBlank(producerTraceId) + && producerTraceId.equals(consumerTraceId); DemoMessageResult result = new DemoMessageResult() .setId(message.getId()) .setScene(message.getScene()) .setContent(message.getContent()) .setSendTime(message.getSendTime()) .setConsumeTime(now) - .setConsumer(consumer); + .setConsumer(consumer) + .setProducerTraceId(producerTraceId) + .setConsumerTraceId(consumerTraceId) + .setTracePropagated(propagated); if (message.getSendTime() != null) { long cost = now.toInstant().toEpochMilli() - message.getSendTime().toInstant().toEpochMilli(); // 处理时钟回拨导致的负值