12 KiB
MQ 消息体序列化结构变更检测 — 方案
版本:v0.4
日期:2026-07-15
状态:Phase M1 + M2 已实现(生产侧 MQ0105 / MQ-K01K04 + 读侧 MQ-R)
关联:复用serialization-schema-checker的 Schema 提取、Diff、企微通知与 CI 框架
业务样本仓:jnpf-java-cloud(RocketMQ + Kafka)
1. 背景与目标
1.1 为什么要做
业务同时使用 RocketMQ 与 Kafka 投递业务对象:
| 中间件 | 典型写法 | 序列化要点 |
|---|---|---|
| RocketMQ | rocketMQTemplate.syncSend(topic:tag, dto) |
Spring MessageConverter(多为 Jackson)把对象变成消息体 |
| Kafka | `kafkaTemplate.send(topic, vo | List)` |
| Kafka 消费 | @KafkaListener + parseObject(message, Xxx.class) |
常以 String 接收后再 Fastjson 反序列化 |
当消息 DTO / VO 删字段、改类型、加包装层时:
- Topic / 重试队列里仍可能有旧结构消息
- 新消费代码反序列化失败,或字段为空导致静默逻辑错误
这与 Redis 缓存「残留旧 value」同一类问题:
| Redis | RocketMQ | Kafka | |
|---|---|---|---|
| 残留形态 | 未过期 key | Topic 积压 / 重试 | Topic 积压 / 消费 lag |
| 路由标识 | key 模式 | topic:tag | topic(一般无 tag;动态后缀可归一 *) |
| 典型序列化 | Fastjson 字符串或 Template 直写 | MessageConverter | Kafka JsonSerializer / 手写 JSON 字符串 |
1.2 目标
在 push 时静态分析 消息体类型的序列化 Schema 是否相对对比区间发生变更,覆盖 RocketMQ + Kafka,并复用现有企微通知 / notify|block 能力。
1.3 非目标
- 不连接真实 Broker,不拉取积压消息做运行时校验
- 不解析依赖 jar 内消息类型(仅本仓
src/main/java) - 不替代权限、幂等、消费失败重试等业务正确性检查
- 不扫仅 Admin 建 Topic 的工具类(如
KafkaTopicUtil,无业务 body) - RabbitMQ 等若后续出现再扩展(当前仓以 RocketMQ / Kafka 为主)
2. 业务调研结论(jnpf-java-cloud)
2.1 RocketMQ
| 写法 | 出现情况 | 策略 |
|---|---|---|
rocketMQTemplate.syncSend(dest, dto) |
高(如钱包扣费) | 纳入 |
asyncSend / syncSendOrderly 等 |
中 | 纳入 |
convertAndSend |
视封装而定 | 纳入 |
MessageBuilder.withPayload 再 send |
中(如 IM 延时消息) | 纳入(MQ04) |
先 JSON.toJSONString 再发 String |
较低 | unwrap 后取类型(MQ05) |
只发 String / byte[] / 无泛型 Message |
有 | 默认忽略 |
样本(资金钱包):
rocketMQTemplate.syncSend(CapitalMqConstants.TOPIC + ":" + tag, req); // WalletDeductReq
@RocketMQMessageListener(...)
public class WalletDeductConsumer implements RocketMQListener<WalletDeductReq> { ... }
2.2 Kafka(已确认需纳入)
| 写法 | 出现情况 | 策略 |
|---|---|---|
kafkaTemplate.send(topic, dto) |
中(值班食安项等) | 纳入 |
kafkaTemplate.send(topic, List<Xxx>) |
中(巡店食安项列表) | 纳入(rootArray) |
kafkaTemplate.send(ProducerRecord) |
中 | 纳入(MQ-K03) |
先 JSON 再 send(topic, json) |
较低 | unwrap(MQ-K04) |
@KafkaListener + String + JSONObject.parseObject(..., Xxx.class) |
有(数据分析中差评) | 读侧补强 MQ-R |
KafkaTopicUtil 仅创建 Topic |
有(租户) | 忽略(无消息体) |
生产样本(巡店):
List<CheckItemDetailVo> thousandsData = ...;
kafkaTemplate.send(topicBuilder.patrolStoreTopic(tenantId), thousandsData);
生产样本(值班):
KafkaTemplate<String, Object> kafkaTemplate;
kafkaTemplate.send(topic, data); // CheckItemDetailVO
消费样本(数据分析):
@KafkaListener(topics = "ftb-evaluate-real-notification${...}", groupId = "...")
public void handleMessage(String message) {
AddedMessageNotificationToVO vo = JSONObject.parseObject(message, AddedMessageNotificationToVO.class);
}
动态 Topic(如按租户拼接)静态推断结果形如 patrol-store-topic:*,与 Redis key * 规则一致。
2.3 「Key」等价物(destination)
| 中间件 | 聚合键形态 | 来源 |
|---|---|---|
| RocketMQ | topic:tag |
字面量、常量、TOPIC + ":" + TAG |
| Kafka | topic |
字面量、常量、topicBuilder.xxx(tenantId) → 前缀+* |
无法解析时:展示表达式 + <font color="comment">(destination 未解析)</font>。
2.4 读侧补强(类比 W06)
| 中间件 | 补强来源 |
|---|---|
| RocketMQ | RocketMQListener<T>、onMessage(T) + @RocketMQMessageListener |
| Kafka | @KafkaListener 方法参数类型;或方法内 parseObject/parseArray(..., Xxx.class)(与现有 W06 共享解析能力) |
开关:detection.mq_read_hints_enabled(默认 true)。
3. 方案总览
3.1 产品形态
并入现有 serialization-schema-checker:
- 同一 CLI / 流水线 / Schema Diff / 企微模板
- 配置增加
mq_patterns(含 RocketMQ + Kafka) - 通知按 Topic / destination 分块;文案统一
Topic -->
3.2 与现有链路
Git Diff → 变更 Java 文件
├─ Redis:W01~W05 + W06 ← 已有
└─ MQ:RocketMQ(MQ01~05)+ Kafka(MQ-K01~K04)
+ Listener / parse 补强(MQ-R) ← 已实现
↓
同一套 TypeSchema / SchemaDiffer / Skeleton / WeCom
对比区间:push before → after。
3.3 核心原则
- 只关心消息体对象 Schema,不关心 Broker / ACL / 限流
- 有结构变更即告警;
block与缓存共用 - 静态分析;仅本仓
src/main/java - RocketMQ 与 Kafka 同一 Diff / 通知模型,仅投递 AST 模式不同
4. 检测模式设计
4.1 生产侧 — RocketMQ
| 模式 ID | 匹配表达式 | 提取 |
|---|---|---|
| MQ01 | rocketMQTemplate.syncSend(dest, payload, …) |
dest、payload 类型 |
| MQ02 | asyncSend / syncSendOrderly / sendOneWay 等 |
同上 |
| MQ03 | convertAndSend(dest, payload) |
同上 |
| MQ04 | MessageBuilder.withPayload(obj) 再 send,或 Message<T> |
payload / T |
| MQ05 | 先 toJSONString/getObjectToString 再 send String |
unwrap 后类型 |
4.2 生产侧 — Kafka
| 模式 ID | 匹配表达式 | 提取 |
|---|---|---|
| MQ-K01 | kafkaTemplate.send(topic, payload) |
topic、payload 类型 |
| MQ-K02 | kafkaTemplate.send(topic, key, payload) |
同上(忽略分区 key) |
| MQ-K03 | send(ProducerRecord) |
topic + value 类型 |
| MQ-K04 | 先 JSON 序列化为 String 再 send(topic, json) |
unwrap 后类型 |
Payload 为 List<Xxx> / Collection 时标记 rootArray,骨架为 JSON 数组(与 Redis List 一致)。
忽略:
- payload 为字面量、纯无结构
String/byte[](无业务类型时) - 仅 Topic Admin API(
AdminClient.createTopics等) - destination 命中
ignore.mq_destinations
4.3 消费侧辅助(不单独告警)
| 模式 ID | 匹配 | 作用 |
|---|---|---|
| MQ-R01 | RocketMQListener<T> / @RocketMQMessageListener |
补强同 destination 生产点 |
| MQ-R02 | @KafkaListener + 参数类型 T(非 String) |
补强同 topic |
| MQ-R03 | Listener 内 parseObject/parseArray(..., Xxx.class) |
补强(复用 parse AST) |
开关:detection.mq_read_hints_enabled(默认 true)。
4.4 Schema Diff
复用现有变更类型与注解规则。
序列化方言:首版按字段名;Jackson / Fastjson / Kafka JsonSerializer 差异必要时用 manual_mappings。
5. 报告与通知
5.1 企微块(RocketMQ / Kafka 统一)
- Topic --> `capital-topic:WALLET_DEDUCT`
> **通道**: `RocketMQ`
> **位置**: `WalletDeductProducer#send:38`
> **类型**: `WalletDeductReq`
> **value值由:** “{...}”
> **变更为:** “{...}”
- Topic --> `patrol-store-food-safe:*`
> **通道**: `Kafka`
> **位置**: `PatrolServiceImpl#sendFoodSafeData:3140`
> **类型**: `List<CheckItemDetailVo>`
> **value值由:** “[{...}]”
> **变更为:** “[{...}]”
- 删除字段橙
warning;新增绿info - destination 未解析时灰色提示
「通道」字段用于区分中间件;若模板求简,可省略通道仅靠 Topic 形态区分。
5.2 CI 控制台
字段明细可标注 RocketMQ / Kafka;再输出与企微一致的 Markdown。
6. 配置草案
detection:
patterns: [W01, W02, W03, W04, W05]
read_hints_enabled: true
mq_patterns:
# RocketMQ
- MQ01
- MQ02
- MQ03
- MQ04
- MQ05
# Kafka
- MQ-K01
- MQ-K02
- MQ-K03
- MQ-K04
mq_read_hints_enabled: true
ignore:
mq_destinations:
- "*:TEST"
- "benchmark:*"
manual_mappings:
- id: wallet-deduct-mq
writer_method: "jnpf.capital.module.wallet.mq.WalletDeductProducer#send"
key_pattern: "capital-topic:WALLET_DEDUCT"
value_type: "jnpf.model.capital.dto.WalletDeductReq"
- id: patrol-kafka-food-safe
writer_method: "jnpf.service.impl.PatrolServiceImpl#sendFoodSafeData"
key_pattern: "*-patrol-store-*" # 按实际 topic 规则调整
value_type: "jnpf.model.analyses.CheckItemDetailVo"
7. 分阶段交付
Phase M1 — MVP(RocketMQ + Kafka 基础投递)✅
| 任务 | 说明 |
|---|---|
| MQ01/MQ02 | RocketMQ syncSend / asyncSend |
| MQ-K01/MQ-K02 | Kafka send(topic, payload) / 三参 send |
| destination 推断 | 字面量、常量、拼接;Kafka 动态 topic → * |
| Schema Diff + 骨架通知 | 复用 ReportBuilder;「通道」行 |
| 夹具 | fixtures/mq/rocket-wallet/、fixtures/mq/kafka-patrol/ |
| 配置 | mq_patterns(含 MQ-K*)、ignore.mq_destinations |
| List 根数组骨架 | List<CheckItemDetailVo> 等(与 M2 需求合并交付) |
验收:
- 删
WalletDeductReq字段 → 企微出现 RocketMQ Topic 骨架变更 - 删
CheckItemDetailVo字段 → 企微出现 Kafka Topic 骨架变更
Phase M2 — 增强 ✅
| 任务 | 说明 |
|---|---|
| MQ03~MQ05、MQ-K03/K04 | convertAndSend、MessageBuilder/Message<T>、JSON 字符串发送、ProducerRecord |
| MQ-R01~R03 | RocketMQ Listener + Kafka @KafkaListener / parse 补强 |
| 夹具 | fixtures/mq/rocket-im/、fixtures/mq/kafka-record/ |
验收:
MessageBuilder.withPayload(DutyImNotice)+asyncSend→ 命中 MQ04,改 VO 字段可告警kafkaTemplate.send(ProducerRecord)→ 命中 MQ-K03- Listener /
parseObject可补强同 Topic 弱类型生产点
8. 风险与限制
| 风险 | 缓解 |
|---|---|
| Kafka topic 按租户动态拼接 | 归一 prefix:*;manual_mappings |
| RocketMQ / Kafka 混用同一 VO | 各投递点独立告警(符合预期) |
| Listener 收 String、parse 在方法深处 | MQ-R03 + parse AST |
| 生产/消费跨模块对不齐 | 同仓索引 + destination 对齐;失败则仅写侧 |
| Jackson / Fastjson / Kafka JsonSerializer 细节差 | 首版字段名;必要时方言或 mapping |
| 只改消费未改生产类型 | 不告警(工具职责是消息体 Schema) |
include_modules 过窄 |
观察期扩大模块或置空全仓 |
9. 决策对齐
| 项 | 结论 |
|---|---|
| 中间件范围 | RocketMQ + Kafka(本仓已确认);Rabbit 暂不纳入 |
| 对比区间 | gitea.event.before → gitea.sha |
| 阻断 | 与现网 mode 共用 |
| 级别 | 不引入 P0/P1/P2 产品展示 |
| 交付 | 同一 jar;patterns 区分 Redis / MQ(含 MQ-K*) |
10. 下一步
按 Phase M1 开发已完成(MQ01/02 + MQ-K01/K02)按 Phase M2 开发已完成(MQ03~05、MQ-K03/K04、MQ-R)- 业务仓验收:
MessageBuilderIM Topic、巡店/值班 Kafka、capital-topic:WALLET_DEDUCT - 观察期按需扩大
include_modules