From f2d85a072ac482941a44ddd952faf63a9460e91b Mon Sep 17 00:00:00 2001 From: lili Date: Thu, 2 Jul 2026 20:04:18 -0700 Subject: [PATCH] =?UTF-8?q?feat(telemetry):=20=E6=89=B9=E9=87=8F=E4=BA=8B?= =?UTF-8?q?=E4=BB=B6=E5=85=A5=E5=8F=A3=E6=94=B9=20RocketMQ=20=E5=BC=82?= =?UTF-8?q?=E6=AD=A5=E6=B6=88=E8=B4=B9(=E5=B9=82=E7=AD=89=C2=B7=E8=81=9A?= =?UTF-8?q?=E5=90=88/=E5=9B=9E=E7=81=8C=E8=A1=8C=E4=B8=BA=E4=B8=8D?= =?UTF-8?q?=E5=8F=98)=20(W-REAL=20R1)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 上报入口 /events/batch·/perf/beacon 从请求线程内同步写改为「轻量校验 + 投递 telemetry-event 主题、 消费侧异步落库+聚合+quality_score 回灌 feed」。灰度开关 telemetry.mq.enabled(默认关=同步回落、现行行为逐字不变)。 - TelemetryEventProducer:逐条投递 event/perf tag、keys=traceId、syncSend 带超时+RocketMQ 自带重试; 发送失败不外抛、按 rejected 交前端重传兜底(不 500 阻断上报)。 - TelemetryEventConsumer:@RocketMQMessageListener 落地骨架,consumeThreadNumber=consumeThreadMax=8 (min=max 破 consumeThreadMax 死参数坑,受 DB 池默认 10 约束);瞬时故障外抛触发 ≤16 次重投转 DLQ, 毒消息吞掉 ack 防重投风暴;幂等由 eventId(uk_event_id)保证重投不双计。 - EventIngestServiceImpl:入口路由(MQ 在席=投递、缺席=同步回落),新增 consumeFromMq 单条事务入口 逐字复用同步写核心 ingestOne——聚合行为与 quality_score 回灌 feed 行为与改造前一致,仅执行位置搬到消费线程; ObjectProvider 软注入生产者 + @Lazy 自代理保回落单条事务真生效。 验证:模块 59 单测绿(既有 EventIngestServiceImplTest 13 个逐字不变通过=回灌行为不变;新增生产者3/消费者4/ MQ路径5);连 mini-infra 真 broker(100.64.0.8:9876)往返 IT 绿(produce→broker→consume→consumeFromMq 转调, 往返序列化保真)。整条 HTTP POST→MQ→真库聚合行 的单体 e2e 待后端全栈窗口(同阶段〇 C1/C2/C3 处置)。 Co-Authored-By: Claude Opus 4.8 (1M context) --- .../game-module-telemetry-server/pom.xml | 7 + .../mq/consumer/TelemetryEventConsumer.java | 93 ++++++--- .../mq/message/TelemetryEventMessage.java | 43 ++++ .../mq/producer/TelemetryEventProducer.java | 96 +++++++++ .../service/event/EventIngestService.java | 12 ++ .../service/event/EventIngestServiceImpl.java | 175 +++++++++++----- .../telemetry/mq/TelemetryMqRoundTripIT.java | 140 +++++++++++++ .../consumer/TelemetryEventConsumerTest.java | 81 ++++++++ .../producer/TelemetryEventProducerTest.java | 94 +++++++++ .../EventIngestServiceImplMqPathTest.java | 188 ++++++++++++++++++ .../src/main/resources/application.yaml | 7 + 11 files changed, 856 insertions(+), 80 deletions(-) create mode 100644 game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/message/TelemetryEventMessage.java create mode 100644 game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/producer/TelemetryEventProducer.java create mode 100644 game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/TelemetryMqRoundTripIT.java create mode 100644 game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/consumer/TelemetryEventConsumerTest.java create mode 100644 game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/producer/TelemetryEventProducerTest.java create mode 100644 game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestServiceImplMqPathTest.java diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/pom.xml b/game-cloud/game-module-telemetry/game-module-telemetry-server/pom.xml index cef2deeb..9a30db59 100644 --- a/game-cloud/game-module-telemetry/game-module-telemetry-server/pom.xml +++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/pom.xml @@ -100,6 +100,13 @@ huijing-spring-boot-starter-rpc + + + org.apache.rocketmq + rocketmq-spring-boot-starter + + com.wanxiang diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/consumer/TelemetryEventConsumer.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/consumer/TelemetryEventConsumer.java index 79b942cb..2ca841c6 100644 --- a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/consumer/TelemetryEventConsumer.java +++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/consumer/TelemetryEventConsumer.java @@ -1,51 +1,80 @@ package com.wanxiang.huijing.game.module.telemetry.mq.consumer; +import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO; +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 lombok.extern.slf4j.Slf4j; +import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; +import org.apache.rocketmq.spring.core.RocketMQListener; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Component; +import jakarta.annotation.Resource; + /** - * 遥测事件 MQ 消费者(脊柱对接点骨架;落原始事件 + 增量聚合 + 回灌) + * 遥测事件 MQ 消费者(切片一 阶段〇 R1)——异步消费事件信封,落原始事件 + 增量聚合 + quality_score 回灌 feed。 * - * 本类为外部依赖对接点,骨架阶段不实现具体消费逻辑(不 mock、不实现),仅声明对接契约与 TODO,保证集成可追溯。 - * 契约见 contracts/api-schemas/telemetry.yaml#x-mq-contracts.consumer: - * - group=telemetry-event-consumer,topic=telemetry-event,无序并发消费(聚合按 gameId 行级原子累加避免覆盖)。 - * - 幂等键 = (traceId, event, ts),命中 game_telemetry_event.uk_dedup 即跳过(INSERT IGNORE 语义), - * 保证至少一次投递下不重复聚合。 - * - 默认重试 ≤16 次仍失败转 DLQ(telemetry-event-dlq),由 T-TEL-19 数据质量监控告警。 + *

脊柱对接点落地(契约 {@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 { + + /** 事件摄取 Service(消费副作用逐字复用同步写核心,保证聚合/回灌行为不变)。 */ + @Resource + private EventIngestService eventIngestService; /** - * TODO 对接点(MQ 消费):补全 RocketMQ 监听注解与消费方法。 - * 待 huijing-spring-boot-starter-mq 的 RocketMQ listener 就绪后,按如下骨架实现: + * 消费一条事件信封:单条事务落库 + 聚合 + quality_score 回灌(幂等)。 * - * @RocketMQMessageListener(topic = "telemetry-event", consumerGroup = "telemetry-event-consumer") - * public class Xxx implements RocketMQListener { - * public void onMessage(EnvelopeMessage msg) { - * // 1) 幂等落原始表(uk_dedup 去重,INSERT IGNORE) - * // 2) 处理状态机 0待聚合 → 1已聚合(EventProcessStatusEnum 校验合法流转) - * // 3) 增量聚合 game_telemetry_game_stat(行级原子累加) - * // 4) 计算 quality_score(公式 TODO) - * // 5) 事务提交后异步回灌 feed/project(Feign 调用不放事务内) - * } - * } - * - * 骨架阶段仅记 INFO 占位说明对接点存在,避免遗漏。 + * @param message MQ 消息体(信封 + 上报时点 userId,见 {@link TelemetryEventMessage}) */ - public TelemetryEventConsumer() { - log.info("[TelemetryEventConsumer] 骨架对接点已就位;MQ 消费/聚合/回灌 待对接(见类注释 TODO)"); + @Override + public void onMessage(TelemetryEventMessage message) { + // 毒消息守卫:结构性坏数据(消息/信封/eventId 为空)重投无益 → 吞掉正常返回(ack),只告警,防无限重投 + if (message == null || message.getEnvelope() == null) { + log.error("[telemetry-mq] 消息或信封为空,丢弃不重投 message={}", message); + return; + } + EnvelopeReqVO envelope = message.getEnvelope(); + try { + // 复用同步写核心(单条 @Transactional):落原始事件(eventId 幂等)→ 增量聚合 → 算 quality_score → 回灌 feed + eventIngestService.consumeFromMq(envelope, message.getUserId()); + log.info("[telemetry-mq] 消费事件完成 eventId={} event={} traceId={}", + envelope.getEventId(), envelope.getEvent(), envelope.getTraceId()); + } catch (Exception e) { + // 瞬时故障(DB 抖动等):外抛让 RocketMQ 重投(单条事务已回滚,幂等保证重试不双计),≤16 次仍失败转 DLQ 告警。 + // 错误路径必须留痕(创始人铁律)。 + log.error("[telemetry-mq] 消费事件失败,外抛触发重投(幂等保证不双计;≤16 次失败转 DLQ)eventId={} event={} traceId={}", + envelope.getEventId(), envelope.getEvent(), envelope.getTraceId(), e); + throw e; + } } - } diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/message/TelemetryEventMessage.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/message/TelemetryEventMessage.java new file mode 100644 index 00000000..7b070584 --- /dev/null +++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/mq/message/TelemetryEventMessage.java @@ -0,0 +1,43 @@ +package com.wanxiang.huijing.game.module.telemetry.mq.message; + +import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO; +import lombok.Data; + +import java.io.Serializable; + +/** + * 遥测事件 MQ 消息体(切片一 阶段〇 R1:批量事件入口改 MQ 异步消费)。 + * + *

契约对齐 {@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 builder = MessageBuilder.withPayload(payload); + if (StringUtils.hasText(envelope.getTraceId())) { + builder.setHeader(RocketMQHeaders.KEYS, envelope.getTraceId()); + } + Message message = builder.build(); + + SendResult result = rocketMQTemplate.syncSend(destination, message, SEND_TIMEOUT_MS); + boolean ok = result != null && result.getSendStatus() == SendStatus.SEND_OK; + if (ok) { + log.info("[telemetry-mq] 投递事件信号 topic={} tag={} eventId={} traceId={} msgId={}", + TOPIC, tag, envelope.getEventId(), envelope.getTraceId(), + result.getMsgId()); + } else { + // 发送非 SEND_OK(如刷盘/主从超时):不外抛,返 false 交调用方按 rejected 处置(客户端重传兜底) + log.warn("[telemetry-mq] 投递事件非 SEND_OK topic={} tag={} eventId={} status={}", + TOPIC, tag, envelope.getEventId(), result == null ? null : result.getSendStatus()); + } + return ok; + } catch (Exception e) { + // 外部交互错误路径必须留痕:MQ 抖动/超时→不外抛(不 500 阻断上报),返 false 计 rejected,前端 SDK 分批重传兜底 + log.warn("[telemetry-mq] 投递事件失败(不外抛,按 rejected 处置由前端重传兜底)topic={} tag={} eventId={} traceId={}", + TOPIC, tag, envelope.getEventId(), envelope.getTraceId(), e); + return false; + } + } +} diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestService.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestService.java index e751cccd..0ea12e4f 100644 --- a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestService.java +++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/main/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestService.java @@ -32,4 +32,16 @@ public interface EventIngestService { */ Boolean ingestBeacon(EnvelopeReqVO reqVO, Long userId); + /** + * MQ 消费侧单条摄取(切片一 阶段〇 R1)——异步消费入口,单条事务落库 + 聚合 + quality_score 回灌。 + * + *

由 {@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 telemetryEventProducerProvider; + /** + * 自身代理引用(切片一 阶段〇 R1):MQ 关闭回落同步写时,逐条经代理调 {@link #consumeFromMq},保证 {@code @Transactional} + * 单条事务真生效(Spring 自调用不走代理、事务失效,故显式经代理)。{@code @Lazy} 破本 Bean 自引用初始化环。 + *

纯 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 received = new ArrayBlockingQueue<>(64); + ObjectMapper om = new ObjectMapper(); + + // ── 真 push consumer:订阅 telemetry-event,收到即把 body 塞队列 ── + DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(group); + consumer.setNamesrvAddr(NAMESRV); + consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET); + consumer.subscribe(TelemetryEventProducer.TOPIC, TelemetryEventProducer.TAG_EVENT); + consumer.registerMessageListener((MessageListenerConcurrently) (msgs, ctx) -> { + msgs.forEach(m -> received.offer(m.getBody())); + return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; + }); + consumer.start(); + // 新消费组首启需等首次 rebalance + 路由拉取,并 establish CONSUME_FROM_LAST_OFFSET 的尾偏移; + // 在此之前投递会被「跳到尾」而漏收(新组 offset 竞态)。等 8s 让其 settle 后再投。 + Thread.sleep(8000); + + // ── 真 producer:构造 RocketMQTemplate(rocketmq-client DefaultMQProducer)→ 真 TelemetryEventProducer ── + org.apache.rocketmq.client.producer.DefaultMQProducer mqProducer = + new org.apache.rocketmq.client.producer.DefaultMQProducer("telemetry-event-producer-e2e"); + mqProducer.setNamesrvAddr(NAMESRV); + mqProducer.start(); + RocketMQTemplate template = new RocketMQTemplate(); + template.setProducer(mqProducer); + template.setMessageConverter(new RocketMQMessageConverter().getMessageConverter()); + TelemetryEventProducer producer = new TelemetryEventProducer(template); + + try { + EnvelopeReqVO env = new EnvelopeReqVO(); + env.setEventId(eventId); + env.setEvent("game_play_end"); + env.setSchemaVersion("v1"); + env.setTs(System.currentTimeMillis()); + env.setTraceId(traceId); + EnvelopeReqVO.ContextVO ctx = new EnvelopeReqVO.ContextVO(); + ctx.setGameId("1024"); + ctx.setVersionId("2048"); + env.setContext(ctx); + + // 1) 真投递到真 broker(返 true 即 broker 已确认 SEND_OK) + boolean sent = producer.publish(env, 99L, TelemetryEventProducer.TAG_EVENT); + assertTrue(sent, "真 broker 应确认 SEND_OK"); + + // 2) 真 broker 回投 → 扫描出我这条(按 eventId 匹配,容忍 topic 上的其它并发流量);≤20s + TelemetryEventMessage roundTrip = null; + long deadline = System.currentTimeMillis() + 20000; + while (System.currentTimeMillis() < deadline) { + byte[] body = received.poll(2, TimeUnit.SECONDS); + if (body == null) { + continue; + } + TelemetryEventMessage m = om.readValue(new String(body, StandardCharsets.UTF_8), TelemetryEventMessage.class); + if (m.getEnvelope() != null && eventId.equals(m.getEnvelope().getEventId())) { + roundTrip = m; + break; + } + } + assertNotNull(roundTrip, "应从真 broker 消费到刚投递的消息(按 eventId 匹配)"); + + // 3) 往返序列化保真:信封逐字段一致 + assertNotNull(roundTrip.getEnvelope(), "消息信封应往返完整"); + assertEquals(eventId, roundTrip.getEnvelope().getEventId(), "eventId 往返一致(幂等真身)"); + assertEquals("game_play_end", roundTrip.getEnvelope().getEvent()); + assertEquals("1024", roundTrip.getEnvelope().getContext().getGameId(), "gameId 往返一致(聚合落点)"); + assertEquals(99L, roundTrip.getUserId(), "userId 随消息携带、往返一致"); + + // 4) 真消费者类处理真 broker 消息 → 转调 consumeFromMq(聚合入口,逐字复用 ingestOne) + TelemetryEventConsumer realConsumer = new TelemetryEventConsumer(); + EventIngestService ingestMock = mock(EventIngestService.class); + ReflectionTestUtils.setField(realConsumer, "eventIngestService", ingestMock); + realConsumer.onMessage(roundTrip); + ArgumentCaptor envCaptor = ArgumentCaptor.forClass(EnvelopeReqVO.class); + verify(ingestMock, timeout(2000)).consumeFromMq(envCaptor.capture(), eq(99L)); + assertEquals(eventId, envCaptor.getValue().getEventId(), "消费者应把真 broker 信封转调聚合入口"); + } finally { + consumer.shutdown(); + mqProducer.shutdown(); + } + } +} diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/consumer/TelemetryEventConsumerTest.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/consumer/TelemetryEventConsumerTest.java new file mode 100644 index 00000000..6cb2d184 --- /dev/null +++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/mq/consumer/TelemetryEventConsumerTest.java @@ -0,0 +1,81 @@ +package com.wanxiang.huijing.game.module.telemetry.mq.consumer; + +import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO; +import com.wanxiang.huijing.game.module.telemetry.mq.message.TelemetryEventMessage; +import com.wanxiang.huijing.game.module.telemetry.service.event.EventIngestService; +import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest; +import org.junit.jupiter.api.Test; +import org.mockito.InjectMocks; +import org.mockito.Mock; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +/** + * {@link TelemetryEventConsumer} 单元测试(纯 Mockito,不触真 MQ)——把守 R1 消费者转调 + 毒消息守卫 + 重试语义。 + * + *

覆盖:① 正常消费转调 {@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> msgCaptor = ArgumentCaptor.forClass(Message.class); + verify(rocketMQTemplate).syncSend(eq("telemetry-event:event"), msgCaptor.capture(), anyLong()); + Object keys = msgCaptor.getValue().getHeaders().get("KEYS"); + assertTrue("trace-abc".equals(keys), "MQ keys 应置 traceId 便于对账"); + } + + /** best-effort(核心断言):syncSend 抛异常 → publish 吞异常不外抛、返 false(前端重传兜底,不 500 阻断上报)。 */ + @Test + void testPublish_bestEffort_swallowsExceptionReturnsFalse() { + doThrow(new RuntimeException("MQ down")).when(rocketMQTemplate) + .syncSend(org.mockito.ArgumentMatchers.anyString(), + (Message) org.mockito.ArgumentMatchers.any(), anyLong()); + + boolean sent = assertDoesNotThrow(() -> producer.publish(envelope(), null, TelemetryEventProducer.TAG_EVENT)); + assertFalse(sent, "发送异常应返回 false(按 rejected 处置)"); + } + + /** 非 SEND_OK(如刷盘超时)→ 返 false(不外抛)。 */ + @Test + @SuppressWarnings("unchecked") + void testPublish_nonSendOk_returnsFalse() { + SendResult flushTimeout = new SendResult(); + flushTimeout.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT); + when(rocketMQTemplate.syncSend(org.mockito.ArgumentMatchers.anyString(), + (Message) org.mockito.ArgumentMatchers.any(), anyLong())).thenReturn(flushTimeout); + + boolean sent = producer.publish(envelope(), 1L, TelemetryEventProducer.TAG_PERF); + assertFalse(sent, "非 SEND_OK 应返回 false"); + } +} diff --git a/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestServiceImplMqPathTest.java b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestServiceImplMqPathTest.java new file mode 100644 index 00000000..6efe5072 --- /dev/null +++ b/game-cloud/game-module-telemetry/game-module-telemetry-server/src/test/java/com/wanxiang/huijing/game/module/telemetry/service/event/EventIngestServiceImplMqPathTest.java @@ -0,0 +1,188 @@ +package com.wanxiang.huijing.game.module.telemetry.service.event; + +import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EnvelopeReqVO; +import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EventBatchReqVO; +import com.wanxiang.huijing.game.module.telemetry.controller.app.event.vo.EventBatchResultVO; +import com.wanxiang.huijing.game.module.telemetry.dal.dataobject.stat.GameStatDO; +import com.wanxiang.huijing.game.module.telemetry.dal.mysql.event.TelemetryEventMapper; +import com.wanxiang.huijing.game.module.telemetry.dal.mysql.stat.GameStatMapper; +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; +import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest; +import org.junit.jupiter.api.Test; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.springframework.beans.factory.ObjectProvider; + +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +/** + * {@link EventIngestServiceImpl} MQ 异步路径单测(切片一 阶段〇 R1)——把守「生产者在席=投递不同步写」路由 + * 与「消费入口 consumeFromMq 逐字复用 ingestOne 聚合/回灌」。与 {@code EventIngestServiceImplTest}(同步回落路径、 + * 既有断言逐字不变)互补,共同坐实「聚合行为与 quality_score 回灌行为不变、仅执行位置从请求线程搬到消费线程」。 + * + * @author 绘境AI + */ +class EventIngestServiceImplMqPathTest extends BaseMockitoUnitTest { + + @InjectMocks + private EventIngestServiceImpl eventIngestService; + + @Mock + private TelemetryEventMapper telemetryEventMapper; + @Mock + private GameStatMapper gameStatMapper; + @Mock + private FeedApi feedApi; + @Mock + private AnomalyFilterService anomalyFilterService; + @Mock + private ObjectProvider telemetryEventProducerProvider; // MQ 生产者软注入 + @Mock + private TelemetryEventProducer telemetryEventProducer; + + /** MQ 开启:ingestBatch 逐条投递、不在请求线程做同步写(DB/feed 零交互)。 */ + @Test + void testIngestBatch_mqEnabled_producesNotSync() { + when(telemetryEventProducerProvider.getIfAvailable()).thenReturn(telemetryEventProducer); + when(telemetryEventProducer.publish(any(), any(), eq(TelemetryEventProducer.TAG_EVENT))).thenReturn(true); + + EventBatchReqVO reqVO = new EventBatchReqVO(); + List batch = new ArrayList<>(); + batch.add(validEnvelope("game_play_start", "t1")); // 合法 → 投递 + batch.add(validEnvelope("like", "t2")); // 合法 → 投递 + batch.add(validEnvelope("not_registered_evt", "t3")); // 未登记 → rejected(不投递) + reqVO.setBatch(batch); + + EventBatchResultVO result = eventIngestService.ingestBatch(reqVO, 99L); + + assertEquals(2, result.getAccepted()); + assertEquals(1, result.getRejected()); + assertEquals("t1", result.getTraceId()); + // 逐条投递 event tag 两次 + verify(telemetryEventProducer, times(2)).publish(any(), eq(99L), eq(TelemetryEventProducer.TAG_EVENT)); + // 关键:请求线程内不做同步写(落库/聚合/回灌全在消费侧) + verifyNoInteractions(telemetryEventMapper); + verifyNoInteractions(gameStatMapper); + verifyNoInteractions(feedApi); + } + + /** MQ 开启但投递失败:该条计 rejected(前端重传兜底),不误报 accepted。 */ + @Test + void testIngestBatch_mqEnabled_publishFailCountedRejected() { + when(telemetryEventProducerProvider.getIfAvailable()).thenReturn(telemetryEventProducer); + when(telemetryEventProducer.publish(any(), any(), anyString())).thenReturn(false); // 投递恒失败 + + EventBatchReqVO reqVO = new EventBatchReqVO(); + List batch = new ArrayList<>(); + batch.add(validEnvelope("game_play_start", "t1")); + reqVO.setBatch(batch); + + EventBatchResultVO result = eventIngestService.ingestBatch(reqVO, 99L); + + assertEquals(0, result.getAccepted()); + assertEquals(1, result.getRejected(), "投递失败应计 rejected 交前端重传"); + } + + /** MQ 开启:beacon 投递 perf tag、恒返 true、不同步写。 */ + @Test + void testIngestBeacon_mqEnabled_producesPerfTag() { + when(telemetryEventProducerProvider.getIfAvailable()).thenReturn(telemetryEventProducer); + when(telemetryEventProducer.publish(any(), any(), eq(TelemetryEventProducer.TAG_PERF))).thenReturn(true); + + Boolean ok = eventIngestService.ingestBeacon(validEnvelope("perf_first_screen", "t-beacon"), 99L); + + assertTrue(ok); + verify(telemetryEventProducer).publish(any(), eq(99L), eq(TelemetryEventProducer.TAG_PERF)); + verifyNoInteractions(telemetryEventMapper); + } + + /** + * 消费入口 consumeFromMq 逐字复用 ingestOne:game_play_end(completed) → 原子累加 → 重读算分(60.00) → 回灌 feed。 + * 坐实「聚合/回灌行为不变」——与 EventIngestServiceImplTest#testIngestOne_playEndTriggersQualityRefresh 同口径, + * 只是入口从 ingestBatch 换成 consumeFromMq(消费线程入口)。 + */ + @Test + void testConsumeFromMq_reusesIngestOne_aggregatesAndBackfills() { + Map props = new HashMap<>(); + props.put("completed", true); + props.put("duration_ms", 5000); + EnvelopeReqVO env = gameEnvelope("game_play_end", "t1", "1024", "2048", props); + when(telemetryEventMapper.selectByEventId(env.getEventId())).thenReturn(null); + GameStatDO accumulated = new GameStatDO(); + accumulated.setId(7L); + accumulated.setPlayCount(0L); + accumulated.setPlayEndCount(1L); + accumulated.setCompletedCount(1L); + accumulated.setTotalDurationMs(5000L); + accumulated.setLikeCount(0L); + accumulated.setShareCount(0L); + when(gameStatMapper.selectByGameAndDate(eq(1024L), any(LocalDate.class))).thenReturn(accumulated); + + eventIngestService.consumeFromMq(env, 99L); + + // 原子累加(play_end=1, completed=1, duration=5000) + verify(gameStatMapper).insertOrAccumulate(eq(1024L), any(LocalDate.class), + eq(0L), eq(1L), eq(1L), eq(5000L), eq(0L), eq(0L)); + // 算分写回 + 回灌 feed(QUALITY_REFRESH,60.0) + verify(gameStatMapper).updateById(any(GameStatDO.class)); + verify(feedApi).upsertRank(any(FeedRankUpsertReqDTO.class)); + // 幂等预查发生(重放同 eventId 不双计的第一道网) + verify(telemetryEventMapper).selectByEventId(env.getEventId()); + } + + /** consumeFromMq 幂等:同 eventId 已落库 → 跳过累加/回灌(重复投递不双计)。 */ + @Test + void testConsumeFromMq_idempotentSkipOnDuplicate() { + EnvelopeReqVO env = gameEnvelope("game_play_start", "t1", "1024", "2048", null); + when(telemetryEventMapper.selectByEventId(env.getEventId())) + .thenReturn(new com.wanxiang.huijing.game.module.telemetry.dal.dataobject.event.TelemetryEventDO()); + + eventIngestService.consumeFromMq(env, 99L); + + verify(telemetryEventMapper, never()).insert(any(com.wanxiang.huijing.game.module.telemetry.dal.dataobject.event.TelemetryEventDO.class)); + verifyNoInteractions(gameStatMapper); + verifyNoInteractions(feedApi); + } + + // ============================== 测试夹具 ============================== + + private static EnvelopeReqVO validEnvelope(String event, String traceId) { + EnvelopeReqVO envelope = new EnvelopeReqVO(); + envelope.setEventId(UUID.randomUUID().toString()); + envelope.setEvent(event); + envelope.setSchemaVersion("v1"); + envelope.setTs(1733600000000L); + envelope.setTraceId(traceId); + return envelope; + } + + private static EnvelopeReqVO gameEnvelope(String event, String traceId, String gameId, String versionId, Map props) { + EnvelopeReqVO envelope = validEnvelope(event, traceId); + EnvelopeReqVO.ContextVO ctx = new EnvelopeReqVO.ContextVO(); + ctx.setGameId(gameId); + ctx.setVersionId(versionId); + envelope.setContext(ctx); + envelope.setProps(props); + return envelope; + } +} diff --git a/game-cloud/huijing-server/src/main/resources/application.yaml b/game-cloud/huijing-server/src/main/resources/application.yaml index a83364bc..004d34f9 100644 --- a/game-cloud/huijing-server/src/main/resources/application.yaml +++ b/game-cloud/huijing-server/src/main/resources/application.yaml @@ -190,6 +190,13 @@ rocketmq: producer: group: ${spring.application.name}_PRODUCER # 生产者分组 +# telemetry 遥测异步(切片一 阶段〇 R1):批量事件入口改 MQ 异步消费的灰度开关 +telemetry: + mq: + # true=上报入口快速 ACK + 投递 telemetry-event 主题、消费侧异步聚合(幂等); + # false(默认)=生产者/消费者 Bean 缺席、回落原请求线程内同步写(现行行为逐字不变)。经 .env 注入按环境开启。 + enabled: ${TELEMETRY_MQ_ENABLED:false} + spring: # Kafka 配置项,对应 KafkaProperties 配置类 kafka: