Files
schemaCheck/docs/MQ序列化结构检测方案.md
dongzi e88408e9b6
All checks were successful
序列化结构检查 / serialization-schema-check (push) Successful in 1s
feat: 一阶段
2026-07-15 16:48:46 +08:00

339 lines
12 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# MQ 消息体序列化结构变更检测 — 方案
> 版本v0.3
> 日期2026-07-15
> 状态:**Phase M1 已实现**MQ01/MQ02 + MQ-K01/MQ-K02M2/M3 未启动
> 关联:复用 `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)` | Spring `KafkaTemplate` + value serializer多为 Json编码对象 |
| 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` | 有 | **默认忽略** |
样本(资金钱包):
```java
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 | 有(租户) | **忽略**(无消息体) |
生产样本(巡店):
```java
List<CheckItemDetailVo> thousandsData = ...;
kafkaTemplate.send(topicBuilder.patrolStoreTopic(tenantId), thousandsData);
```
生产样本(值班):
```java
KafkaTemplate<String, Object> kafkaTemplate;
kafkaTemplate.send(topic, data); // CheckItemDetailVO
```
消费样本(数据分析):
```java
@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 产品形态
并入现有 `serialization-schema-checker`
- 同一 CLI / 流水线 / Schema Diff / 企微模板
- 配置增加 `mq_patterns`(含 RocketMQ + Kafka
- 通知按 **Topic / destination** 分块;文案统一 `Topic -->`
### 3.2 与现有链路
```text
Git Diff → 变更 Java 文件
├─ RedisW01~W05 + W06 ← 已有
└─ MQRocketMQMQ01~+ KafkaMQ-K*
+ Listener / parse 补强MQ-R ← 本方案
同一套 TypeSchema / SchemaDiffer / Skeleton / WeCom
```
对比区间push **`before``after`**。
### 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 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)` | 补强(可复用 W06 检测器) |
开关:`detection.mq_read_hints_enabled`(默认 true
### 4.4 Schema Diff
复用现有变更类型与注解规则。
序列化方言首版按字段名Jackson / Fastjson / Kafka JsonSerializer 差异必要时用 `manual_mappings`
---
## 5. 报告与通知
### 5.1 企微块RocketMQ / Kafka 统一)
```markdown
- 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. 配置草案
```yaml
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.before``gitea.sha` |
| 阻断 | 与现网 `mode` 共用 |
| 级别 | 不引入 P0/P1/P2 产品展示 |
| 交付 | 同一 jarpatterns 区分 Redis / MQ含 MQ-K* |
---
## 10. 下一步
1. ~~按 Phase M1 开发~~ **已完成**MQ01/02 + MQ-K01/K02
2. Phase M2MQ03~05、MQ-K03/K04、MQ-R 读侧补强
3. 业务仓验收 Topic`capital-topic:WALLET_DEDUCT`RocketMQ、巡店/值班 Kafka topic
4. 观察期按需扩大 `include_modules`