Files
schemaCheck/docs/MQ序列化结构检测方案.md
2026-07-15 11:29:12 +08:00

12 KiB
Raw Blame History

MQ 消息体序列化结构变更检测 — 方案

版本v0.2
日期2026-07-15
状态:方案已落地(含 Kafka开发未启动
关联:复用 cache-schema-checker 的 Schema 提取、Diff、企微通知与 CI 框架
业务样本仓:jnpf-java-cloudRocketMQ + Kafka


1. 背景与目标

1.1 为什么要做

业务同时使用 RocketMQKafka 投递业务对象:

中间件 典型写法 序列化要点
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 视封装而定 纳入
JSON.toJSONString 再发 String 较低 unwrap 后取类型
只发 String / byte[] / MessageExt 默认忽略

样本(资金钱包):

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
@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 共享解析能力)

3. 方案总览

3.1 产品形态

并入现有 cache-schema-checker(或后续更名 serialization-schema-checker

  • 同一 CLI / 流水线 / Schema Diff / 企微模板
  • 配置增加 mq_patterns(含 RocketMQ + Kafka
  • 通知按 Topic / destination 分块;文案统一 Topic -->

3.2 与现有链路

Git Diff → 变更 Java 文件
    ├─ RedisW01~W05 + W06                    ← 已有
    └─ MQRocketMQMQ01~+ KafkaMQ-K*
         + Listener / parse 补强MQ-R         ← 本方案
              ↓
        同一套 TypeSchema / SchemaDiffer / Skeleton / WeCom

对比区间push beforeafter

3.3 核心原则

  1. 只关心消息体对象 Schema,不关心 Broker / ACL / 限流
  2. 有结构变更即告警block 与缓存共用
  3. 静态分析;仅本仓 src/main/java
  4. 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 payload 类型
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) / ListenableFuture 封装若可解析 topic + value 类型
MQ-K04 先 JSON 序列化为 String 再 send(topic, json) unwrap 后类型

Payload 为 List<Xxx> / Collection 时标记 rootArray,骨架为 JSON 数组(与 Redis List 一致)。

忽略

  • payload 为字面量、纯无结构 String/byte[](无业务类型时)
  • 仅 Topic Admin APIAdminClient.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) 补强(可复用 W06 检测器)

开关: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 — MVPRocketMQ + 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

验收

  1. WalletDeductReq 字段 → 企微出现 RocketMQ Topic 骨架变更
  2. CheckItemDetailVo 字段 → 企微出现 Kafka Topic 骨架变更

Phase M2 — 增强

任务 说明
MQ03~MQ05、MQ-K03/K04 convertAndSend、MessageBuilder、JSON 字符串发送、ProducerRecord
MQ-R01~R03 RocketMQ Listener + Kafka @KafkaListener / parse 补强
List 根数组骨架 巡店 List<CheckItemDetailVo>

Phase M3 — 运营

任务 说明
积压风险提示文案 统一提示评估消费 lag / 积压
更多夹具 值班 Kafka、中差评 Listener、IM Favorite 等
分 webhook / 标题前缀 缓存 vs MQ 可选拆分

8. 风险与限制

风险 缓解
Kafka topic 按租户动态拼接 归一 prefix:*manual_mappings
RocketMQ / Kafka 混用同一 VO 各投递点独立告警(符合预期)
Listener 收 String、parse 在方法深处 MQ-R03 + 复用 W06 AST
生产/消费跨模块对不齐 同仓索引 + destination 对齐;失败则仅写侧
Jackson / Fastjson / Kafka JsonSerializer 细节差 首版字段名;必要时方言或 mapping
只改消费未改生产类型 不告警(工具职责是消息体 Schema

9. 决策对齐

结论
中间件范围 RocketMQ + Kafka本仓已确认Rabbit 暂不纳入
对比区间 gitea.event.beforegitea.sha
阻断 与现网 mode 共用
级别 不引入 P0/P1/P2 产品展示
交付 同一 jarpatterns 区分 Redis / MQ含 MQ-K*

10. 下一步

  1. 评审本方案RocketMQ + Kafka 模式表)
  2. Phase M1 开发MQ01/02 + MQ-K01/K02
  3. 回写 配置说明.md / CI集成说明.md 正式配置项
  4. 验收 Topic 建议:capital-topic:WALLET_DEDUCTRocketMQ、巡店/值班 Kafka topic