脊柱对接点落地(契约 {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.consumer}): + * topic={@link TelemetryEventProducer#TOPIC}、consumerGroup=telemetry-event-consumer、并发消费(无序)。 + * 消费副作用逐字复用同步写核心 {@code EventIngestServiceImpl#ingestOne}(经 {@link EventIngestService#consumeFromMq} + * 单条事务包裹)——聚合行为与 quality_score 回灌 feed 行为与改造前逐字一致,只是执行位置从请求线程搬到消费线程。 * - * 消费副作用(待实现): - * 1. 落 game_telemetry_event(原始事件,幂等)—— 经 TelemetryEventMapper insert,依赖 uk_dedup 去重。 - * 2. 增量聚合 game_telemetry_game_stat(play/like/quality 等按 gameId+statDate 行级累加,唯一约束 uk_game_date 保证幂等)。 - * 3. quality_score 计算(运营质量分 0-100;输入:avgDurationMs/completedCount/reportCount/loadFailCount 等)—— - * TODO 对接点:计算公式待产品/数据侧确定后落地,不在骨架阶段臆造公式。 - * 4. 回灌:play_count/like_count → project(game_project);quality_score → feed 排序信号—— - * TODO 对接点:待 project/feed 的 -api(Feign)就绪后在事务提交后异步回灌(远程调用不放事务内)。 + *
幂等:以信封 {@code eventId}(uk_event_id)为幂等键——同步写核心先 selectByEventId 预查、命中即跳过累加, + * 叠加 uk_event_id 唯一约束硬后盾;「至少一次投递」下重复消息不重复聚合(契约 V10 起 uk_dedup 已被 uk_event_id 取代)。 + * + *
并发上限({@code consumeThreadMax} 死参数坑):rocketmq-client 无界消费队列致 {@code consumeThreadMax} + * 不单独生效(线程数永不超 core),故与 {@code consumeThreadNumber} 同设为 8(min=max)让消费并发上限真为 8; + * MQ 只 cap 投递(消费)速率,不是「在跑生成并发」那种业务并发。8 取值受 DB 连接池(默认 10)约束留头寸, + * 每条消费=一次快 DB 写事务,可按吞吐与池大小调(改这里需同调池)。 + * + *
重试与 DLQ(外部交互纪律):消费抛异常则 RocketMQ 重投(默认 ≤16 次)仍失败转 DLQ(%RETRY%/%DLQ% 组), + * 由 T-TEL-19 数据质量监控告警。故对瞬时故障(DB 抖动)本方法外抛让其重试(单条事务已回滚、幂等保证重试不双计); + * 对毒消息(信封为空等结构性坏数据,重投无益)吞掉正常返回(ack),只告警留痕,防无限重投风暴。 + * + *
装配条件:{@code telemetry.mq.enabled=true} 才装配(与生产者同一灰度开关)——默认关时不起消费者、不连 NameServer。
*
* @author 绘境AI
*/
@Slf4j
@Component
-public class TelemetryEventConsumer {
+@ConditionalOnProperty(prefix = "telemetry.mq", name = "enabled", havingValue = "true")
+@RocketMQMessageListener(
+ topic = TelemetryEventProducer.TOPIC,
+ consumerGroup = "telemetry-event-consumer",
+ consumeThreadNumber = 8, // min=max=8:消费并发上限真为 8(受 DB 池默认 10 约束留头寸;见类注释死参数坑)
+ consumeThreadMax = 8 // rocketmq-client 无界队列致本参不单独生效,须与 consumeThreadNumber 同值
+)
+public class TelemetryEventConsumer implements RocketMQListener 契约对齐 {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.producer.messageBody}:核心是
+ * 逐条投递的事件信封 {@link EnvelopeReqVO};此外携带上报时点解析出的登录 {@code userId}(传输关切、
+ * 非契约字段)——异步消费脱离了 Web 请求线程,无法再从 SecurityContext 取登录态,故须在生产时随信封带上,
+ * 保证消费侧算 user_id 兜底与审计列填充与同步写路径逐字一致(见 {@code EventIngestServiceImpl#consumeFromMq})。
+ *
+ * 幂等真身 = 信封 {@code eventId}(消费侧命中 {@code uk_event_id} 去重);MQ 消息 keys 用 traceId 仅便于排障对账。
+ * 经 rocketmq-spring 默认 Jackson 转换器序列化/反序列化往返;字段均为 POJO,无循环引用。
+ *
+ * @author 绘境AI
+ */
+@Data
+public class TelemetryEventMessage implements Serializable {
+
+ private static final long serialVersionUID = 1L;
+
+ /** 单条事件信封(幂等真身 eventId 在其内;消费侧据此落库 + 聚合 + 回灌)。 */
+ private EnvelopeReqVO envelope;
+
+ /**
+ * 上报时点解析的登录用户 ID(匿名为 null)。异步消费无 Web 请求线程、取不到登录态,故随消息携带,
+ * 供消费侧 resolveUserId 兜底 + 匿名审计列兜底判定,与同步写路径口径一致。
+ */
+ private Long userId;
+
+ /** 空构造(Jackson 反序列化需要)。 */
+ public TelemetryEventMessage() {
+ }
+
+ public TelemetryEventMessage(EnvelopeReqVO envelope, Long userId) {
+ this.envelope = envelope;
+ this.userId = userId;
+ }
+}
diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/producer/TelemetryEventProducer.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/producer/TelemetryEventProducer.java
new file mode 100644
index 00000000..bfccea17
--- /dev/null
+++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/producer/TelemetryEventProducer.java
@@ -0,0 +1,96 @@
+package com.wanxiang.huijing.game.module.telemetry.mq.producer;
+
+import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
+import com.wanxiang.huijing.game.module.telemetry.mq.message.TelemetryEventMessage;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.rocketmq.client.producer.SendResult;
+import org.apache.rocketmq.client.producer.SendStatus;
+import org.apache.rocketmq.spring.core.RocketMQTemplate;
+import org.apache.rocketmq.spring.support.RocketMQHeaders;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
+import org.springframework.stereotype.Component;
+import org.springframework.util.StringUtils;
+
+/**
+ * 遥测事件 MQ 生产者(切片一 阶段〇 R1)——上报通道轻量校验后逐条投递事件信封到 {@code telemetry-event} 主题,
+ * 把「落库 + 聚合 + quality_score 回灌」从请求线程内同步写改为消费侧异步完成(幂等)。
+ *
+ * 契约对齐 {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.producer}:topic={@link #TOPIC}、
+ * tag=event|perf({@code /events/batch} 用 event、{@code /perf/beacon} 用 perf)、keys=traceId、逐条投递、orderly=false。
+ *
+ * 外部交互纪律(创始人铁律):syncSend 带发送超时 {@link #SEND_TIMEOUT_MS};RocketMQ 生产者自带
+ * 发送失败重试(retryTimesWhenSendFailed 默认 2)作瞬时抖动兜底;幂等由消费侧 {@code eventId}(uk_event_id)去重
+ * 保证「至少一次投递」不重复聚合。发送最终失败不外抛——返回 {@code false} 交由调用方按「受理失败(rejected)」处置,
+ * 由前端 SDK 分批重传作最外层兜底(遥测=分析数据、可容忍偶发丢失且客户端重试,绝不 500 阻断上报)。
+ *
+ * 装配条件:{@code telemetry.mq.enabled=true} 才装配(与消费者同一灰度开关)——默认关时本 Bean 缺席,
+ * {@code EventIngestServiceImpl} 经 {@code ObjectProvider} 取空、回落原同步写路径(现行行为逐字不变)。
+ *
+ * @author 绘境AI
+ */
+@Slf4j
+@Component
+@ConditionalOnProperty(prefix = "telemetry.mq", name = "enabled", havingValue = "true")
+public class TelemetryEventProducer {
+
+ /** 遥测事件主题(与消费者 {@code TelemetryEventConsumer} 同名订阅;契约 telemetry-event)。 */
+ public static final String TOPIC = "telemetry-event";
+
+ /** 事件通道 tag(/events/batch)。 */
+ public static final String TAG_EVENT = "event";
+
+ /** 性能兜底通道 tag(/perf/beacon)。 */
+ public static final String TAG_PERF = "perf";
+
+ /** 发送超时(毫秒)——外部交互必设超时,避免上报请求被 MQ 抖动无限拖住。 */
+ private static final long SEND_TIMEOUT_MS = 3000L;
+
+ /** RocketMQ 模板(rocketmq-spring 自动装配;name-server 见 application*.yaml → mini-infra 100.64.0.8:9876)。 */
+ private final RocketMQTemplate rocketMQTemplate;
+
+ public TelemetryEventProducer(RocketMQTemplate rocketMQTemplate) {
+ this.rocketMQTemplate = rocketMQTemplate;
+ }
+
+ /**
+ * 投递一条事件信封到 {@link #TOPIC}(同步发 + 超时;成功 SEND_OK 返 true,否则/异常返 false)。
+ *
+ * @param envelope 已通过轻量校验的单条信封
+ * @param userId 上报时点登录用户 ID(匿名为 null;随消息携带供消费侧口径一致)
+ * @param tag 消息 tag({@link #TAG_EVENT} / {@link #TAG_PERF})
+ * @return 是否投递成功(true=broker 已确认 SEND_OK;false=发送失败/异常,调用方按 rejected 处置)
+ */
+ public boolean publish(EnvelopeReqVO envelope, Long userId, String tag) {
+ // 目的地格式 topic:tag(rocketmq-spring 约定)
+ String destination = TOPIC + ":" + tag;
+ try {
+ // 消息体 = 信封 + userId;keys 置 traceId(便于全链路排障对账,非幂等键——幂等真身是 eventId)
+ TelemetryEventMessage payload = new TelemetryEventMessage(envelope, userId);
+ MessageBuilder 由 {@code TelemetryEventConsumer} 调用(消费线程);MQ 关闭时也作 {@code ingestBatch}/{@code ingestBeacon}
+ * 的同步回落逐条入口(经代理保证单条事务)。逐字复用同步写核心,保证聚合/回灌行为与改造前一致;
+ * 幂等由信封 eventId(uk_event_id)保证,重复投递不重复聚合。
+ *
+ * @param reqVO 单条信封(已在生产侧通过轻量校验)
+ * @param userId 上报时点登录用户 ID(匿名为 null;随 MQ 消息携带,保证消费侧口径与同步写一致)
+ */
+ void consumeFromMq(EnvelopeReqVO reqVO, Long userId);
+
}
diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestServiceImpl.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestServiceImpl.java
index a66287d9..bdf903ab 100644
--- a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestServiceImpl.java
+++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestServiceImpl.java
@@ -9,6 +9,7 @@ import com.wanxiang.huijing.game.module.telemetry.dal.mysql.event.TelemetryEvent
import com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat.GameStatMapper;
import com.wanxiang.huijing.game.module.telemetry.enums.EventProcessStatusEnum;
import com.wanxiang.huijing.game.module.telemetry.enums.TelemetryEventEnum;
+import com.wanxiang.huijing.game.module.telemetry.mq.producer.TelemetryEventProducer;
import com.wanxiang.huijing.game.module.telemetry.service.quality.AnomalyFilterService;
import com.wanxiang.huijing.game.module.feed.api.FeedApi;
import com.wanxiang.huijing.game.module.feed.dto.FeedRankUpsertReqDTO;
@@ -18,6 +19,8 @@ import com.wanxiang.huijing.framework.common.util.servlet.ServletUtils;
import com.wanxiang.huijing.framework.security.core.LoginUser;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.ObjectProvider;
+import org.springframework.context.annotation.Lazy;
import org.springframework.security.authentication.UsernamePasswordAuthenticationToken;
import org.springframework.security.core.context.SecurityContextHolder;
import org.springframework.stereotype.Service;
@@ -99,61 +102,53 @@ public class EventIngestServiceImpl implements EventIngestService {
@Resource
private AnomalyFilterService anomalyFilterService;
- /*
- * TODO 对接点(MQ 生产,§3.5 仅留契约+TODO 不实现):拆"上报快速 ACK + 异步消费聚合"时,
- * 注入 RocketMQ 生产者投递 topic=telemetry-event(幂等键=eventId,对应 uk_event_id),消费侧复用本类同步写聚合逻辑。
- * 见 contracts/api-schemas/telemetry.yaml#x-mq-contracts。MVP 同步写已满足闭环,不引入 MQ。
+ /**
+ * 遥测事件 MQ 生产者软注入(切片一 阶段〇 R1):{@code telemetry.mq.enabled=true} 才在席——上报入口把
+ * 「落库 + 聚合 + quality_score 回灌」改为快速 ACK + 投递 {@code telemetry-event} 主题、消费侧异步完成(幂等)。
+ * {@code enabled=false}(默认)时缺席、{@link ObjectProvider#getIfAvailable()} 取空 → 回落原同步写路径
+ * (现行行为逐字不变,既有单测据此仍绿)。纯 Mockito 单测不注入本 Provider,{@link #resolveProducer()} 空值守卫回落同步写。
*/
+ @Resource
+ private ObjectProvider 纯 Mockito 单测无 Spring 代理,本字段为 null,{@link #invokeConsumeFromMq}/{@link #ingestBatchSyncFallback}
+ * 空值守卫退化为直调 {@code this}(事务注解在单测本就无 Spring 语义、不影响 mapper 交互断言)。
+ */
+ @Lazy
+ @Resource
+ private EventIngestService self;
+
+ // ============================== 上报入口(MQ 异步路由 / 同步回落)==============================
+
+ /**
+ * 批量摄取(切片一 阶段〇 R1 路由):MQ 开启则逐条轻量校验后投递 {@code telemetry-event}(快速 ACK、消费侧异步聚合);
+ * MQ 关闭则回落原同步写(逐条单事务经代理)。批量上限超限整体拒绝——先于任何投递/落库,与改造前口径一致。
+ */
@Override
- @Transactional(rollbackFor = Exception.class) // §3.5 C5 同步写:落库+聚合+回灌+置已聚合在同一本地事务,任一步失败整批回滚
public EventBatchResultVO ingestBatch(EventBatchReqVO reqVO, Long userId) {
- // 2026-06-10 e2e 二轮实证(P0-1 补口):匿名上报是多级联写(事件 insert → stat 原子累加 → feed rank 回灌
- // → markAggregated UPDATE),其中 markAggregated 的 updateById 运行时实证仍把 updater=null 写进
- // UPDATE SET(与 M-b 三因叠加同形态)。逐写点 set 兜底易漏级联(feed 模块回灌同线程同雷),故在入口
- // 对匿名请求注入系统身份(照 AigcGenerateExecutor.executeWithSystemIdentity 范本的 Web 线程版),
- // 使全链 DefaultDBFieldHandler 填充 creator/updater="0";finally 清理防 Web 线程上下文泄漏。
- boolean injected = injectSystemIdentityIfAnonymous(userId);
- try {
- // 批量上限:超限整体拒绝(统一错误码出口,前端据此分批重传)
- if (reqVO.getBatch().size() > MAX_BATCH_SIZE) {
- throw exception(TELEMETRY_BATCH_SIZE_EXCEEDED);
- }
- // 本次批量受理回执 traceId(便于排障;优先复用首条信封 traceId,缺失则生成)
- String batchTraceId = resolveBatchTraceId(reqVO);
-
- int accepted = 0;
- int rejected = 0;
- // 部分成功语义:逐条轻量校验,通过则同步写(落库+聚合+回灌),失败计 rejected(不中断整批)
- for (EnvelopeReqVO envelope : reqVO.getBatch()) {
- if (!isEnvelopeValid(envelope)) {
- rejected++;
- continue;
- }
- ingestOne(envelope, userId);
- accepted++;
- }
- // 关键链路日志:受理成功记 INFO,含回执 traceId 与计数,便于全链路排障
- log.info("[ingestBatch] 批量同步写完成 batchTraceId={}, accepted={}, rejected={}", batchTraceId, accepted, rejected);
-
- EventBatchResultVO result = new EventBatchResultVO();
- result.setAccepted(accepted);
- result.setRejected(rejected);
- result.setTraceId(batchTraceId);
- return result;
- } finally {
- if (injected) {
- SecurityContextHolder.clearContext();
- }
+ // 批量上限:超限整体拒绝(统一错误码出口,前端据此分批重传)——先判,任何投递/落库前
+ if (reqVO.getBatch().size() > MAX_BATCH_SIZE) {
+ throw exception(TELEMETRY_BATCH_SIZE_EXCEEDED);
}
+ TelemetryEventProducer producer = resolveProducer();
+ if (producer != null) {
+ // MQ 异步路径:逐条校验 + 投递 event tag,快速 ACK 返回 accepted/rejected(不在请求线程做聚合)
+ return dispatchBatchToMq(reqVO, userId, producer);
+ }
+ // 同步回落(MQ 关闭 / 单测):逐条单事务落库 + 聚合 + 回灌,行为与改造前一致
+ return ingestBatchSyncFallback(reqVO, userId);
}
+ /**
+ * 兜底单条上报(切片一 阶段〇 R1 路由):校验失败恒返 true(不阻塞前端卸载);校验通过 MQ 开启则投递 perf tag、
+ * 关闭则同步写。fire-and-forget,任何路径恒返 true。
+ */
@Override
- @Transactional(rollbackFor = Exception.class) // 兜底通道单条同步写:与 ingestBatch 同口径(落库+聚合+回灌原子)
public Boolean ingestBeacon(EnvelopeReqVO reqVO, Long userId) {
- // 兜底通道:页面卸载场景,校验失败不抛错(返回受理成功,不触发前端异常);
- // 校验通过则同步写——若同步写抛错(极少:DB 故障),由全局异常处理返回错误响应,
- // 但 sendBeacon 为 fire-and-forget 不读响应、不阻塞卸载(事务该回滚仍回滚,保证不脏写)。
+ // 兜底通道:页面卸载场景,校验失败不抛错(返回受理成功,不触发前端异常)
if (!isEnvelopeValid(reqVO)) {
log.warn("[ingestBeacon] 信封校验未通过,丢弃 eventId={}, event={}, traceId={}",
reqVO == null ? null : reqVO.getEventId(),
@@ -161,11 +156,30 @@ public class EventIngestServiceImpl implements EventIngestService {
reqVO == null ? null : reqVO.getTraceId());
return Boolean.TRUE;
}
- // 匿名注入系统身份:与 ingestBatch 同口径(级联写审计列兜底),finally 清理
+ TelemetryEventProducer producer = resolveProducer();
+ if (producer != null) {
+ // MQ 异步:投递 perf tag;fire-and-forget 不阻塞卸载,投递失败由前端 sendBeacon 天然不重试(性能埋点可容忍丢失)
+ producer.publish(reqVO, userId, TelemetryEventProducer.TAG_PERF);
+ return Boolean.TRUE;
+ }
+ // 同步回落:单条事务落库 + 聚合 + 回灌(与 ingestBatch 同口径)
+ invokeConsumeFromMq(reqVO, userId);
+ return Boolean.TRUE;
+ }
+
+ /**
+ * MQ 消费侧单条摄取(切片一 阶段〇 R1):消费线程或同步回落逐条入口——单条 {@code @Transactional} 落库 + 聚合 + 回灌,
+ * 逐字复用同步写核心 {@link #ingestOne}(聚合/回灌行为与改造前一致);匿名注入系统身份补审计列,finally 清理。
+ * 2026-06-10 e2e P0-1 补口沿用:匿名上报多级联写(事件 insert / stat 原子累加 / feed rank 回灌 / markAggregated UPDATE)
+ * 的 DefaultDBFieldHandler 审计列填充统一取 "0";消费线程无 Web 请求上下文,ServletUtils.getClientIP() 返 null
+ * 由异常剔除桩按仅 anonId 维度处理(本就 MQ-ready)。
+ */
+ @Override
+ @Transactional(rollbackFor = Exception.class) // 单条同步写:落库+聚合+回灌+置已聚合在同一本地事务,任一步失败该条回滚(幂等保证重试不双计)
+ public void consumeFromMq(EnvelopeReqVO reqVO, Long userId) {
boolean injected = injectSystemIdentityIfAnonymous(userId);
try {
ingestOne(reqVO, userId);
- return Boolean.TRUE;
} finally {
if (injected) {
SecurityContextHolder.clearContext();
@@ -173,6 +187,71 @@ public class EventIngestServiceImpl implements EventIngestService {
}
}
+ // ---------- 路由内部实现 ----------
+
+ /** 解析 MQ 生产者(软注入 + 空值守卫):Provider 缺席/未注入(enabled=false 或单测)→ null → 回落同步写。 */
+ private TelemetryEventProducer resolveProducer() {
+ return telemetryEventProducerProvider == null ? null : telemetryEventProducerProvider.getIfAvailable();
+ }
+
+ /**
+ * MQ 异步路径批量投递:逐条轻量校验,通过则投递 event tag(成功计 accepted、投递失败计 rejected 交前端重传),
+ * 校验失败计 rejected。返回 accepted/rejected + 回执 traceId(部分成功语义、非整批失败)。
+ */
+ private EventBatchResultVO dispatchBatchToMq(EventBatchReqVO reqVO, Long userId, TelemetryEventProducer producer) {
+ String batchTraceId = resolveBatchTraceId(reqVO);
+ int accepted = 0;
+ int rejected = 0;
+ for (EnvelopeReqVO envelope : reqVO.getBatch()) {
+ if (!isEnvelopeValid(envelope)) {
+ rejected++;
+ continue;
+ }
+ // 投递成功=已受理(accepted);投递失败=按受理失败(rejected),前端分批重传兜底(幂等由 eventId 保证不双计)
+ if (producer.publish(envelope, userId, TelemetryEventProducer.TAG_EVENT)) {
+ accepted++;
+ } else {
+ rejected++;
+ }
+ }
+ log.info("[ingestBatch] 批量投递 MQ 完成 batchTraceId={}, accepted={}, rejected={}", batchTraceId, accepted, rejected);
+ EventBatchResultVO result = new EventBatchResultVO();
+ result.setAccepted(accepted);
+ result.setRejected(rejected);
+ result.setTraceId(batchTraceId);
+ return result;
+ }
+
+ /**
+ * 同步回落批量摄取(MQ 关闭 / 单测):逐条轻量校验,通过则经代理调 {@link #consumeFromMq}(单条事务),失败计 rejected。
+ * 与改造前的差异仅事务粒度(整批→逐条,更贴契约「非整批失败」语义、坏事件不连累好事件);聚合/回灌逻辑逐字不变。
+ */
+ private EventBatchResultVO ingestBatchSyncFallback(EventBatchReqVO reqVO, Long userId) {
+ String batchTraceId = resolveBatchTraceId(reqVO);
+ int accepted = 0;
+ int rejected = 0;
+ for (EnvelopeReqVO envelope : reqVO.getBatch()) {
+ if (!isEnvelopeValid(envelope)) {
+ rejected++;
+ continue;
+ }
+ invokeConsumeFromMq(envelope, userId);
+ accepted++;
+ }
+ log.info("[ingestBatch] 批量同步写完成(MQ 关闭回落)batchTraceId={}, accepted={}, rejected={}", batchTraceId, accepted, rejected);
+ EventBatchResultVO result = new EventBatchResultVO();
+ result.setAccepted(accepted);
+ result.setRejected(rejected);
+ result.setTraceId(batchTraceId);
+ return result;
+ }
+
+ /** 经代理调 {@link #consumeFromMq}(保证单条 {@code @Transactional} 真生效);单测无代理时空值守卫直调 this。 */
+ private void invokeConsumeFromMq(EnvelopeReqVO envelope, Long userId) {
+ EventIngestService target = (self != null) ? self : this;
+ target.consumeFromMq(envelope, userId);
+ }
+
/**
* 匿名上报注入系统身份(2026-06-10 e2e 二轮 P0-1 补口,范本=AigcGenerateExecutor.executeWithSystemIdentity)。
*
diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/TelemetryMqRoundTripIT.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/TelemetryMqRoundTripIT.java
new file mode 100644
index 00000000..897088b3
--- /dev/null
+++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/TelemetryMqRoundTripIT.java
@@ -0,0 +1,140 @@
+package com.wanxiang.huijing.game.module.telemetry.mq;
+
+import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
+import com.wanxiang.huijing.game.module.telemetry.mq.consumer.TelemetryEventConsumer;
+import com.wanxiang.huijing.game.module.telemetry.mq.message.TelemetryEventMessage;
+import com.wanxiang.huijing.game.module.telemetry.mq.producer.TelemetryEventProducer;
+import com.wanxiang.huijing.game.module.telemetry.service.event.EventIngestService;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
+import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
+import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
+import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
+import org.apache.rocketmq.spring.core.RocketMQTemplate;
+import org.apache.rocketmq.spring.support.RocketMQMessageConverter;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.condition.EnabledIfSystemProperty;
+import org.mockito.ArgumentCaptor;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.nio.charset.StandardCharsets;
+import java.util.UUID;
+import java.util.concurrent.ArrayBlockingQueue;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+
+/**
+ * R1「MQ 消费聚合真跑」集成证据——连 mini-infra 真 RocketMQ({@code 100.64.0.8:9876} · Tailscale 直连、绕系统代理)
+ * 跑一次 produce → broker → consume → 转调 {@code consumeFromMq} 全链路,坐实生产者/消费者/消息序列化对真 broker 成立。
+ *
+ * 覆盖边界:真 broker 的投递/消费/往返序列化 + 两个生产类({@link TelemetryEventProducer} /
+ * {@link TelemetryEventConsumer})在真 broker 上的行为。聚合落库(consumeFromMq→ingestOne→game_telemetry_game_stat)
+ * 的逻辑由 {@code EventIngestServiceImplMqPathTest#testConsumeFromMq_reusesIngestOne_aggregatesAndBackfills} 单测坐实
+ * (逐字复用 ingestOne,不重造);本 IT 用 mock 的 EventIngestService 捕获转调,不连 MySQL。整条
+ * HTTP POST→MQ→真库聚合行 的单体 e2e 需后端窗口(MySQL/Redis/Nacos 全栈),与阶段〇 C1/C2/C3 集成验证同window处置。
+ *
+ * 默认跳过:{@code @EnabledIfSystemProperty(telemetry.mq.e2e=1)}——常规 {@code mvn test} 不连 broker、不跑本 IT;
+ * 出证据时显式 {@code mvn test -Dtelemetry.mq.e2e=1 -Dtest=TelemetryMqRoundTripIT}。
+ *
+ * @author 绘境AI
+ */
+@EnabledIfSystemProperty(named = "telemetry.mq.e2e", matches = "1")
+class TelemetryMqRoundTripIT {
+
+ /** mini-infra NameServer(Tailscale 直连;无系统代理 env,直连不被 fake-ip 拦)。 */
+ private static final String NAMESRV = System.getProperty("telemetry.mq.namesrv", "100.64.0.8:9876");
+
+ @Test
+ void produce_consume_roundTrip_againstRealBroker() throws Exception {
+ String eventId = "e2e-" + UUID.randomUUID();
+ String traceId = "trace-" + UUID.randomUUID();
+ // 独立测试消费组,earliest 起消费防错过刚发的消息
+ String group = "telemetry-event-consumer-e2e-" + System.currentTimeMillis();
+ BlockingQueue 覆盖:① 正常消费转调 {@code consumeFromMq(envelope,userId)};② 毒消息(消息/信封为空)吞掉不转调、不外抛(ack 防无限重投);
+ * ③ 瞬时故障({@code consumeFromMq} 抛异常)外抛让 RocketMQ 重投(≤16 次转 DLQ,幂等保证不双计)。
+ *
+ * @author 绘境AI
+ */
+class TelemetryEventConsumerTest extends BaseMockitoUnitTest {
+
+ @Mock
+ private EventIngestService eventIngestService;
+
+ @InjectMocks
+ private TelemetryEventConsumer consumer;
+
+ private static EnvelopeReqVO envelope() {
+ EnvelopeReqVO env = new EnvelopeReqVO();
+ env.setEventId("evt-1");
+ env.setEvent("game_play_end");
+ env.setSchemaVersion("v1");
+ env.setTs(1733600000000L);
+ env.setTraceId("trace-abc");
+ return env;
+ }
+
+ /** 正常路:消费转调 consumeFromMq(envelope, userId),逐字复用同步写核心。 */
+ @Test
+ void testOnMessage_dispatchesToConsumeFromMq() {
+ EnvelopeReqVO env = envelope();
+ TelemetryEventMessage msg = new TelemetryEventMessage(env, 99L);
+
+ assertDoesNotThrow(() -> consumer.onMessage(msg));
+
+ verify(eventIngestService).consumeFromMq(eq(env), eq(99L));
+ }
+
+ /** 毒消息守卫①:消息为 null → 吞掉、不转调、不外抛(ack 防无限重投)。 */
+ @Test
+ void testOnMessage_nullMessage_swallowedNoDispatch() {
+ assertDoesNotThrow(() -> consumer.onMessage(null));
+ verify(eventIngestService, never()).consumeFromMq(any(), any());
+ }
+
+ /** 毒消息守卫②:信封为 null → 吞掉、不转调、不外抛。 */
+ @Test
+ void testOnMessage_nullEnvelope_swallowedNoDispatch() {
+ TelemetryEventMessage msg = new TelemetryEventMessage(null, 99L);
+ assertDoesNotThrow(() -> consumer.onMessage(msg));
+ verify(eventIngestService, never()).consumeFromMq(any(), any());
+ }
+
+ /** 瞬时故障:consumeFromMq 抛异常 → onMessage 外抛(触发 RocketMQ 重投;幂等保证重试不双计)。 */
+ @Test
+ void testOnMessage_transientFailure_rethrowsForRetry() {
+ EnvelopeReqVO env = envelope();
+ TelemetryEventMessage msg = new TelemetryEventMessage(env, null);
+ doThrow(new RuntimeException("DB down")).when(eventIngestService).consumeFromMq(eq(env), isNull());
+
+ assertThrows(RuntimeException.class, () -> consumer.onMessage(msg));
+ }
+}
diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/producer/TelemetryEventProducerTest.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/producer/TelemetryEventProducerTest.java
new file mode 100644
index 00000000..c2f5aa66
--- /dev/null
+++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/producer/TelemetryEventProducerTest.java
@@ -0,0 +1,94 @@
+package com.wanxiang.huijing.game.module.telemetry.mq.producer;
+
+import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO;
+import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
+import org.apache.rocketmq.client.producer.SendResult;
+import org.apache.rocketmq.client.producer.SendStatus;
+import org.apache.rocketmq.spring.core.RocketMQTemplate;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.springframework.messaging.Message;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * {@link TelemetryEventProducer} 单元测试(纯 Mockito,不触真 MQ)——把守 R1 生产者投递契约与外部交互纪律。
+ *
+ * 覆盖:① 正常投递到 {@code telemetry-event:event}、返 true;② SEND_OK 判定;③ syncSend 抛异常时
+ * {@code publish} 不外抛、返 false(交调用方按 rejected 处置,前端重传兜底,绝不 500 阻断上报);
+ * ④ 非 SEND_OK 返 false;⑤ keys 置 traceId。
+ *
+ * @author 绘境AI
+ */
+class TelemetryEventProducerTest extends BaseMockitoUnitTest {
+
+ @Mock
+ private RocketMQTemplate rocketMQTemplate;
+
+ @InjectMocks
+ private TelemetryEventProducer producer;
+
+ /** 构造一条最简合法信封(含 eventId/traceId)。 */
+ private static EnvelopeReqVO envelope() {
+ EnvelopeReqVO env = new EnvelopeReqVO();
+ env.setEventId("evt-1");
+ env.setEvent("game_play_start");
+ env.setSchemaVersion("v1");
+ env.setTs(1733600000000L);
+ env.setTraceId("trace-abc");
+ return env;
+ }
+
+ /** 正常路:投递到 topic:tag=telemetry-event:event,SEND_OK → 返 true,keys=traceId。 */
+ @Test
+ @SuppressWarnings("unchecked")
+ void testPublish_sendOk_returnsTrueAndKeysTraceId() {
+ SendResult ok = new SendResult();
+ ok.setSendStatus(SendStatus.SEND_OK);
+ ok.setMsgId("mid-1");
+ when(rocketMQTemplate.syncSend(eq(TelemetryEventProducer.TOPIC + ":" + TelemetryEventProducer.TAG_EVENT),
+ (Message>) org.mockito.ArgumentMatchers.any(), anyLong())).thenReturn(ok);
+
+ boolean sent = producer.publish(envelope(), 99L, TelemetryEventProducer.TAG_EVENT);
+
+ assertTrue(sent, "SEND_OK 应返回 true");
+ // 校验投递目的地 + keys header = traceId
+ ArgumentCaptor