feat(telemetry): 批量事件入口改 RocketMQ 异步消费(幂等·聚合/回灌行为不变) (W-REAL R1)

上报入口 /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) <noreply@anthropic.com>
This commit is contained in:
lili 2026-07-02 20:04:18 -07:00
parent 4d220e77da
commit f2d85a072a
11 changed files with 856 additions and 80 deletions

View File

@ -100,6 +100,13 @@
<artifactId>huijing-spring-boot-starter-rpc</artifactId>
</dependency>
<!-- RocketMQ切片一 阶段〇 R1telemetry 批量事件入口改 MQ 异步消费RocketMQTemplate / @RocketMQMessageListener。
starter-mq 以 optional=true 声明本 starter不传递故本模块须显式直依赖版本由 huijing 依赖管理统一控制。 -->
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>
<!-- 测试 -->
<dependency>
<groupId>com.wanxiang</groupId>

View File

@ -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-consumertopic=telemetry-event无序并发消费聚合按 gameId 行级原子累加避免覆盖
* - 幂等键 = (traceId, event, ts)命中 game_telemetry_event.uk_dedup 即跳过INSERT IGNORE 语义
* 保证至少一次投递下不重复聚合
* - 默认重试 16 次仍失败转 DLQtelemetry-event-dlq T-TEL-19 数据质量监控告警
* <p><b>脊柱对接点落地</b>契约 {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.consumer}
* topic={@link TelemetryEventProducer#TOPIC}consumerGroup=telemetry-event-consumer并发消费无序
* 消费副作用逐字复用同步写核心 {@code EventIngestServiceImpl#ingestOne} {@link EventIngestService#consumeFromMq}
* 单条事务包裹<b>聚合行为与 quality_score 回灌 feed 行为与改造前逐字一致</b>只是执行位置从请求线程搬到消费线程
*
* 消费副作用待实现
* 1. game_telemetry_event原始事件幂等 TelemetryEventMapper insert依赖 uk_dedup 去重
* 2. 增量聚合 game_telemetry_game_statplay/like/quality 等按 gameId+statDate 行级累加唯一约束 uk_game_date 保证幂等
* 3. quality_score 计算运营质量分 0-100输入avgDurationMs/completedCount/reportCount/loadFailCount
* TODO 对接点计算公式待产品/数据侧确定后落地不在骨架阶段臆造公式
* 4. 回灌play_count/like_count projectgame_projectquality_score feed 排序信号
* TODO 对接点 project/feed -apiFeign就绪后在事务提交后异步回灌远程调用不放事务内
* <p><b>幂等</b>以信封 {@code eventId}uk_event_id为幂等键同步写核心先 selectByEventId 预查命中即跳过累加
* 叠加 uk_event_id 唯一约束硬后盾至少一次投递下重复消息不重复聚合契约 V10 uk_dedup 已被 uk_event_id 取代
*
* <p><b>并发上限{@code consumeThreadMax} 死参数坑</b>rocketmq-client 无界消费队列致 {@code consumeThreadMax}
* 不单独生效线程数永不超 core故与 {@code consumeThreadNumber} <b>同设为 8min=max</b>让消费并发上限真为 8
* MQ cap 投递消费速率不是在跑生成并发那种业务并发8 取值受 DB 连接池默认 10约束留头寸
* 每条消费=一次快 DB 写事务可按吞吐与池大小调改这里需同调池
*
* <p><b>重试与 DLQ外部交互纪律</b>消费抛异常则 RocketMQ 重投默认 16 仍失败转 DLQ%RETRY%/%DLQ%
* T-TEL-19 数据质量监控告警故对<b>瞬时故障</b>DB 抖动本方法<b>外抛让其重试</b>单条事务已回滚幂等保证重试不双计
* <b>毒消息</b>信封为空等结构性坏数据重投无益<b>吞掉正常返回</b>ack只告警留痕防无限重投风暴
*
* <p><b>装配条件</b>{@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<TelemetryEventMessage> {
/** 事件摄取 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<EnvelopeMessage> {
* 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/projectFeign 调用不放事务内
* }
* }
*
* 骨架阶段仅记 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 次失败转 DLQeventId={} event={} traceId={}",
envelope.getEventId(), envelope.getEvent(), envelope.getTraceId(), e);
throw e;
}
}
}

View File

@ -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 异步消费
*
* <p>契约对齐 {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.producer.messageBody}核心是
* 逐条投递的事件信封 {@link EnvelopeReqVO}此外携带上报时点解析出的登录 {@code userId}传输关切
* 非契约字段异步消费脱离了 Web 请求线程无法再从 SecurityContext 取登录态故须在生产时随信封带上
* 保证消费侧算 user_id 兜底与审计列填充与同步写路径逐字一致 {@code EventIngestServiceImpl#consumeFromMq}
*
* <p>幂等真身 = 信封 {@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;
}
}

View File

@ -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 回灌从请求线程内同步写改为消费侧异步完成幂等
*
* <p><b>契约对齐</b> {@code contracts/api-schemas/telemetry.yaml#x-mq-contracts.producer}topic={@link #TOPIC}
* tag=event|perf{@code /events/batch} event{@code /perf/beacon} perfkeys=traceId逐条投递orderly=false
*
* <p><b>外部交互纪律创始人铁律</b>syncSend 带发送超时 {@link #SEND_TIMEOUT_MS}RocketMQ 生产者自带
* 发送失败重试retryTimesWhenSendFailed 默认 2作瞬时抖动兜底<b>幂等</b>由消费侧 {@code eventId}uk_event_id去重
* 保证至少一次投递不重复聚合发送最终失败<b>不外抛</b>返回 {@code false} 交由调用方按受理失败(rejected)处置
* 由前端 SDK 分批重传作最外层兜底遥测=分析数据可容忍偶发丢失且客户端重试绝不 500 阻断上报
*
* <p><b>装配条件</b>{@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_OKfalse=发送失败/异常调用方按 rejected 处置
*/
public boolean publish(EnvelopeReqVO envelope, Long userId, String tag) {
// 目的地格式 topic:tagrocketmq-spring 约定
String destination = TOPIC + ":" + tag;
try {
// 消息体 = 信封 + userIdkeys traceId便于全链路排障对账非幂等键幂等真身是 eventId
TelemetryEventMessage payload = new TelemetryEventMessage(envelope, userId);
MessageBuilder<TelemetryEventMessage> builder = MessageBuilder.withPayload(payload);
if (StringUtils.hasText(envelope.getTraceId())) {
builder.setHeader(RocketMQHeaders.KEYS, envelope.getTraceId());
}
Message<TelemetryEventMessage> 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;
}
}
}

View File

@ -32,4 +32,16 @@ public interface EventIngestService {
*/
Boolean ingestBeacon(EnvelopeReqVO reqVO, Long userId);
/**
* MQ 消费侧单条摄取切片一 阶段 R1异步消费入口单条事务落库 + 聚合 + quality_score 回灌
*
* <p> {@code TelemetryEventConsumer} 调用消费线程MQ 关闭时也作 {@code ingestBatch}/{@code ingestBeacon}
* 的同步回落逐条入口经代理保证单条事务逐字复用同步写核心保证聚合/回灌行为与改造前一致
* 幂等由信封 eventIduk_event_id保证重复投递不重复聚合
*
* @param reqVO 单条信封已在生产侧通过轻量校验
* @param userId 上报时点登录用户 ID匿名为 null MQ 消息携带保证消费侧口径与同步写一致
*/
void consumeFromMq(EnvelopeReqVO reqVO, Long userId);
}

View File

@ -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-contractsMVP 同步写已满足闭环不引入 MQ
/**
* 遥测事件 MQ 生产者软注入切片一 阶段 R1{@code telemetry.mq.enabled=true} 才在席上报入口把
* 落库 + 聚合 + quality_score 回灌改为快速 ACK + 投递 {@code telemetry-event} 主题消费侧异步完成幂等
* <p>{@code enabled=false}默认时缺席{@link ObjectProvider#getIfAvailable()} 取空 回落原同步写路径
* 现行行为逐字不变既有单测据此仍绿 Mockito 单测不注入本 Provider{@link #resolveProducer()} 空值守卫回落同步写
*/
@Resource
private ObjectProvider<TelemetryEventProducer> telemetryEventProducerProvider;
/**
* 自身代理引用切片一 阶段 R1MQ 关闭回落同步写时逐条经代理调 {@link #consumeFromMq}保证 {@code @Transactional}
* 单条事务真生效Spring 自调用不走代理事务失效故显式经代理{@code @Lazy} 破本 Bean 自引用初始化环
* <p> 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 tagfire-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 清理
* <p>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
*

View File

@ -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;
/**
* R1MQ 消费聚合真跑集成证据 mini-infra RocketMQ{@code 100.64.0.8:9876} · Tailscale 直连绕系统代理
* 跑一次 produce broker consume 转调 {@code consumeFromMq} 全链路坐实生产者/消费者/消息序列化对真 broker 成立
*
* <p><b>覆盖边界</b> broker 的投递/消费/往返序列化 + 两个生产类{@link TelemetryEventProducer} /
* {@link TelemetryEventConsumer}在真 broker 上的行为聚合落库consumeFromMqingestOnegame_telemetry_game_stat
* 的逻辑由 {@code EventIngestServiceImplMqPathTest#testConsumeFromMq_reusesIngestOne_aggregatesAndBackfills} 单测坐实
* 逐字复用 ingestOne不重造 IT mock EventIngestService 捕获转调不连 MySQL整条
* HTTP POSTMQ真库聚合行 的单体 e2e 需后端窗口MySQL/Redis/Nacos 全栈与阶段 C1/C2/C3 集成验证同window处置
*
* <p><b>默认跳过</b>{@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 NameServerTailscale 直连;无系统代理 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<byte[]> 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构造 RocketMQTemplaterocketmq-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<EnvelopeReqVO> 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();
}
}
}

View File

@ -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 消费者转调 + 毒消息守卫 + 重试语义
*
* <p>覆盖 正常消费转调 {@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));
}
}

View File

@ -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 生产者投递契约与外部交互纪律
*
* <p>覆盖 正常投递到 {@code telemetry-event:event} true SEND_OK 判定 syncSend 抛异常时
* {@code publish} <b>不外抛</b> 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:eventSEND_OK → 返 truekeys=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<Message<?>> 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");
}
}

View File

@ -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<TelemetryEventProducer> 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<EnvelopeReqVO> 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<EnvelopeReqVO> 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 逐字复用 ingestOnegame_play_end(completed) 原子累加 重读算分(60.00) 回灌 feed
* 坐实聚合/回灌行为不变 EventIngestServiceImplTest#testIngestOne_playEndTriggersQualityRefresh 同口径
* 只是入口从 ingestBatch 换成 consumeFromMq消费线程入口
*/
@Test
void testConsumeFromMq_reusesIngestOne_aggregatesAndBackfills() {
Map<String, Object> 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));
// 算分写回 + 回灌 feedQUALITY_REFRESH60.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<String, Object> 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;
}
}

View File

@ -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: