ARB-013 P0 修复:PRODUCER_TOPIC_MAP 从 edu.notification.* 改为 edu.notify.notification.* kafka.consumer 新增 3 考试实时事件消费(exam.extended/force_submitted/question_reordered) 嵌套 payload 解包支持 + topic-map 扩展 新增 6 测试数据文件(docker-notify + 5 kafka 事件 json) 101 单元测试通过 + Docker 真实环境验证
29 KiB
msg 模块 Next Steps v2(上下游依赖 + 状态澄清)
模块:msg(消息中台,HTTP 3007 + gRPC 50056) 更新日期:2026-07-14(v2:ARB-013 落地 + 3 考试实时事件 + 嵌套 payload 解包 + 状态澄清) 状态:13 gRPC RPC + 27 REST endpoints + 16 类 Kafka 消费 + 4 类 Kafka 发布,101 单元测试通过,Docker 本地真实环境验证全部通过
关联:
1. v2 本轮完成的工作
1.1 ARB-013 Kafka topic 命名统一(P0 blocker 修复)
原因:push-gateway nextstep-v2.md §5.2 标注 P0 blocker —— msg Outbox 发布到旧 topic edu.notification.requested,push-gateway 已对齐到 edu.notify.notification.sent,导致 Kafka 链路断裂。
改动:
src/shared/kafka/topic-map.ts—PRODUCER_TOPIC_MAP从edu.notification.*改为edu.notify.notification.*(4 topic + FALLBACK_TOPIC)src/shared/outbox/outbox.publisher.ts— 注释更新src/shared/outbox/outbox.schema.ts— 注释更新
验证结果(Docker 真实 Kafka):
{"eventId":"r4jto...","eventType":"notification.sent","topic":"edu.notify.notification.sent","msg":"Outbox message published"}
Kafka topic 列表确认 edu.notify.notification.sent + edu.notify.notification.read 已创建。
1.2 新增 homework.assigned 消费(core-edu 事件补全)
原因:core-edu 发布 edu.teaching.homework.assigned 事件,但 msg CONSUMER_TOPICS 未包含此 topic。
改动:
src/shared/kafka/topic-map.ts— CONSUMER_TOPICS 添加edu.teaching.homework.assignedsrc/shared/kafka/kafka.consumer.ts— routeEvent 添加 case +handleHomeworkAssigned()方法(向所有 studentIds 发送新作业通知)
1.3 新增 3 个考试实时事件消费(core-edu P3.14)
原因:core-edu nextstep-v2.md §3.2 要求 msg 消费 3 个考试实时事件 topic,生成通知后转发到 edu.notify.notification.*,供 push-gateway 推送到 WebSocket。
事件流:
core-edu → edu.teaching.exam.extended → msg 消费 → edu.notify.notification.sent → push-gateway → WebSocket → student-portal
改动:
src/shared/kafka/topic-map.ts— CONSUMER_TOPICS 添加 3 个 topic(16 类事件总计)src/shared/kafka/kafka.consumer.ts— routeEvent 添加 3 个 case + 3 个 handler +broadcastExamRealtimeEvent()辅助方法
新增 topic:
| Topic | 触发场景 | msg 处理 |
|---|---|---|
edu.teaching.exam.extended |
教师延长考试时间 | 向学生发送"考试时间延长通知" |
edu.teaching.exam.force_submitted |
教师强制提交考试 | 向学生发送"考试强制提交通知" |
edu.teaching.exam.question_reordered |
教师调整题目顺序 | 向学生发送"题目顺序调整通知" |
handler 逻辑:
- 读取
payload.studentIds,为每个学生创建 in_app 通知(type=exam) - ChannelDispatcher 自动调用 push-gateway HTTP
/internal/push(实时推送) - Outbox 自动发布
edu.notify.notification.sent(异步 Kafka 推送) - 若 payload 无 studentIds,记录 warn 并跳过(等待 core-edu 补全 payload)
1.4 嵌套 payload 解包(core-edu Outbox 格式兼容)
原因:core-edu nextstep-v2.md §2.4 明确事件 payload 是嵌套结构 { event_id, event_type, payload: { examId, ... } },但 msg 原有 handler 直接从根取字段(假设扁平结构)。
改动:
src/shared/kafka/kafka.consumer.ts—routeEvent()内部统一解包:若payload.payload存在且为对象,则取payload.payload作为 businessPayload 传给 handler;否则直接用 payload(兼容扁平格式)- 所有 handler 接收 businessPayload(已解包的业务字段),不再处理外层包装
- 考试实时事件 handler 额外接收 rawPayload(用于取
event_id作为 groupId)
兼容性:同时支持 core-edu Outbox 嵌套格式和测试用扁平格式。
1.5 groupId 长度安全(防御性 truncate)
原因:Kafka 消息无 eventId header 时,extractEventId() 降级生成 topic:partition:offset(如 edu.teaching.homework.assigned:0:0 = 35 字符),超过 group_id 字段 varchar(32) 限制,导致 ER_DATA_TOO_LONG。
改动:
handleExamPublished/handleHomeworkAssigned— groupId 优先用payload.eventId(业务事件 ID,cuid2 24 字符);降级路径的 Kafka eventId 若超过 32 字符则 truncatebroadcastExamRealtimeEvent()— 同样从rawPayload.event_id取业务 eventId 作为 groupId,truncate 防御
1.6 sendBatch eventId 幂等过滤(v1 已完成,v2 验证)
v1 改动:
src/notifications/notifications.repository.ts— 新增findExistingEventIds(eventIds)批量查询src/notifications/notifications.service.ts— sendBatch 批量 INSERT 前先过滤已存在的 eventId
v2 Docker 验证:重复 eventId 跳过,DB COUNT=1(幂等生效)。
1.7 pino ESM 导入修复
原因:Docker 构建报 TS2349: This expression is not callable,pino 默认导出在 ESM 模式下不可调用。
改动:src/shared/observability/logger.ts — import pino from 'pino' → import { pino } from 'pino'
1.8 Dockerfile pnpm prune 修复
原因:pnpm prune --prod 触发 husky prepare 脚本(sh: husky: not found)并要求 TTY 确认。
改动:services/msg/Dockerfile — 替换为 pnpm install --prod --ignore-scripts + ENV CI=true HUSKY=0
2. msg 当前能力清单(v2 更新)
2.1 gRPC 13 RPC(端口 50056)
| Service | RPC | 说明 |
|---|---|---|
| NotificationService | SendNotification | 单条发送 |
| NotificationService | ListNotifications | 列表(分页+过滤) |
| NotificationService | MarkAsRead | 标记已读 |
| NotificationService | SearchNotifications | ES 全文检索 |
| NotificationService | RecallNotification | 撤回广播 |
| NotificationService | BatchSendNotification | 批量发送 |
| NotificationService | BatchMarkAsRead | 批量已读 |
| NotificationService | GetUnreadCount | 未读计数 |
| NotificationService | DeleteNotification | 删除通知 |
| NotificationPreferenceService | GetPreferences | 查询偏好 |
| NotificationPreferenceService | UpdatePreferences | 更新偏好 |
| NotificationTemplateService | CreateTemplate/GetTemplate/ListTemplates/UpdateTemplate/DeleteTemplate/RenderTemplate | 模板 CRUD + 渲染 |
2.2 REST endpoints(端口 3007,27 个)
/notifications/*(10 endpoints):send/batch/user/:userId/unread-count/:id/read/batch/read/read-all/search/recall/:id/preferences/user/:userId(2 endpoints):GET/PUT/templates/*(6 endpoints):CRUD + render/announcements/*(9 endpoints):CRUD + publish/archive/pin/read/healthz、/readyz、/metrics
2.3 Kafka 消费(16 类事件)
| 上游 | Topic | msg 处理 |
|---|---|---|
| iam | edu.identity.user.created |
发送欢迎通知 |
| iam | edu.identity.user.updated |
仅幂等标记 |
| iam | edu.identity.user.deleted |
仅幂等标记 |
| iam | edu.identity.user.role_changed |
发送角色变更通知 |
| iam | edu.identity.role.created |
仅幂等标记 |
| iam | edu.identity.role.updated |
向受影响用户发送权限变更通知 |
| core-edu | edu.teaching.exam.published |
向学生发送考试通知 |
| core-edu | edu.teaching.homework.assigned |
向学生发送作业通知 |
| core-edu | edu.teaching.assignment.submitted |
向教师发送提交通知 |
| core-edu | edu.teaching.assignment.graded |
向学生发送批改通知 |
| core-edu | edu.teaching.grade.recorded |
向学生发送成绩通知 |
| core-edu | edu.teaching.attendance.recorded |
向家长发送出勤通知 |
| core-edu | edu.teaching.exam.extended |
向学生发送考试延长通知(P3.14 新增) |
| core-edu | edu.teaching.exam.force_submitted |
向学生发送强制提交通知(P3.14 新增) |
| core-edu | edu.teaching.exam.question_reordered |
向学生发送题目调整通知(P3.14 新增) |
| data-ana | edu.insight.mastery.updated |
向学生发送学情预警通知 |
2.4 Kafka 发布(4 类事件,ARB-013 命名)
edu.notify.notification.sentedu.notify.notification.readedu.notify.notification.recallededu.notify.notification.failed- 兜底:
edu.notify.notification.events
2.5 Docker 验证结果(v2 本轮)
| 测试项 | 结果 |
|---|---|
| 容器健康检查 /healthz | ✅ {"status":"ok"} HTTP 200 |
| KafkaConsumer 订阅 16 topic | ✅ 全部加入 consumer group |
Outbox 发布 topic = edu.notify.notification.sent |
✅ ARB-013 生效 |
| homework.assigned 消费(扁平 payload) | ✅ 2 条通知入库 |
| homework.assigned 消费(降级路径 truncate) | ✅ group_id 截断到 32 字符 |
| exam.extended 消费(嵌套 payload) | ✅ 2 条通知入库,group_id=evt-exam-ext-001 |
| exam.force_submitted 消费(嵌套 payload) | ✅ 1 条通知入库 |
| exam.question_reordered 消费(嵌套 payload) | ✅ 2 条通知入库 |
Outbox 自动转发 5 条 edu.notify.notification.sent |
✅ 全部发布成功 |
| 101 单元测试 | ✅ 全部通过 |
| typecheck + lint | ✅ 零错误 |
3. 上游依赖状态澄清(⚠️ 纠正错误状态)
3.1 ⚠️ teacher-bff nextstep-v2.md §2.5 状态错误
teacher-bff nextstep-v2.md §2.5 标注 3 项 msg 依赖为"⏳ 待 msg 实现",实际全部已实现:
| teacher-bff 标注 | 实际状态 | 说明 |
|---|---|---|
gRPC ListNotifications(userId) ⏳ 待 msg 实现 |
✅ 已实现 | NotificationService.ListNotifications |
gRPC MarkNotificationRead(id) ⏳ 待 msg 实现 |
✅ 已实现 | NotificationService.MarkAsRead |
gRPC ListAnnouncements(filter) ⏳ 待 msg 实现 |
⚠️ REST 已实现,gRPC 未实现 | msg 提供 REST GET /announcements,未提供 gRPC ListAnnouncements |
需 teacher-bff(ai03)更新:
- ListNotifications / MarkNotificationRead 状态改为 ✅
- ListAnnouncements 改用 REST
GET http://msg:3007/announcements(msg 公告为 REST only,ARB-008 限制 RPC 总数 13,不新增 AnnouncementService gRPC) - §5 联调待办 #4"msg gRPC 联调 ⏳"可标记为就绪(msg :50056 已运行)
3.2 ⚠️ parent-bff nextstep-v2.md §4.4 状态错误
parent-bff nextstep-v2.md §4.4 标注 2 项 msg 依赖为"⚠️ proto 未定义 RPC",实际 proto 已定义:
| parent-bff 标注 | 实际状态 | 说明 |
|---|---|---|
getNotificationPreferences(parentId) ⚠️ proto 未定义 |
✅ 已定义 | NotificationPreferenceService.GetPreferences |
updateNotificationPreferences(...) ⚠️ proto 未定义 |
✅ 已定义 | NotificationPreferenceService.UpdatePreferences |
GET /healthz ⏳ msg 服务容器未运行 |
✅ 已运行 | msg :3007/healthz + :50056 gRPC 已就绪 |
需 parent-bff(ai04)更新:
- getNotificationPreferences / updateNotificationPreferences 状态改为 ✅
- RPC 命名对齐:proto 实际为
GetPreferences/UpdatePreferences(非getNotificationPreferences/updateNotificationPreferences),parent-bff gRPC client 需用 proto 标准命名 - /healthz 状态改为 ✅
- §6.2 联调待办 #4"msg 服务容器启动 ⏳"可标记为就绪
- §7"给 msg:补全 NotificationPreferencesService"可标记为已完成
3.3 core-edu 事件 payload 字段确认
core-edu nextstep-v2.md §2.4 确认事件 payload 为嵌套结构:
{
"event_id": "UUID",
"aggregate_id": "examId",
"event_type": "exam.extended",
"occurred_at": 1234567890,
"payload": {
"examId": "...",
"classId": "...",
"subjectId": "...",
"extensionSeconds": 300,
"newDuration": 7200
},
"metadata": { "schema_version": "v1", "trace_id": "...", "user_id": "..." }
}
msg v2 已兼容:routeEvent() 自动解包 payload.payload 到 businessPayload。
⚠️ 待 core-edu 确认:3 个考试实时事件的 payload 是否包含 studentIds 字段?
- msg handler 依赖
payload.studentIds确定通知接收人 - 若 payload 无 studentIds,msg 记录 warn 并跳过(不创建通知)
- core-edu §2.4 示例仅展示 exam.extended 的 examId/classId/subjectId/extensionSeconds/newDuration,未列出 studentIds
- 需 core-edu 确认:考试实时事件 payload 是否应包含
studentIds: string[]
4. 下游依赖状态
4.1 push-gateway(ai09)— ✅ 全部对齐
| 协作项 | 状态 | 说明 |
|---|---|---|
Kafka topic edu.notify.notification.sent |
✅ 已对齐 | msg Outbox 发布到此 topic,push-gateway 消费此 topic |
HTTP POST /internal/push |
✅ 已对齐 | msg ChannelDispatcher 调用 push-gateway 实时推送 |
鉴权头 X-Internal-Key: PUSH_INTERNAL_TOKEN |
✅ 已对齐 | msg push-gateway.client.ts 已使用 |
请求体 { userId, event, data } camelCase |
✅ 已对齐 | msg 已使用 camelCase |
push-gateway nextstep-v2.md §5.2 标注的 P0 blocker 已解除。
4.2 api-gateway(ai01)— ⚠️ 缺 announcements 路由
| 路由 | 状态 | 说明 |
|---|---|---|
/api/v1/notifications/* → msg:3007 |
✅ 已配置 | |
/api/v1/messages/* → msg:3007 |
✅ 已配置 | |
/api/v1/announcements/* → msg:3007 |
⚠️ 未配置 | msg 公告 REST API 在 /announcements/*,需 api-gateway 新增路由 |
需 api-gateway(ai01)新增:/api/v1/announcements/* → msg:3007 路由代理。
4.3 teacher-bff(ai03)— ⚠️ 需接入公告 REST
| 依赖 | 状态 | 说明 |
|---|---|---|
| gRPC ListNotifications | ✅ proto 已定义 | teacher-bff 标注"⏳"有误,实际已实现 |
| gRPC MarkAsRead | ✅ proto 已定义 | teacher-bff 标注"⏳"有误,实际已实现 |
| REST GET /announcements | ✅ msg 已实现 | teacher-bff 期望 gRPC,但 msg 公告为 REST only |
| REST POST /announcements | ✅ msg 已实现 | admin-portal 公告管理需要 |
| REST PUT /announcements/:id/publish | ✅ msg 已实现 | admin-portal 公告发布 |
| REST PUT /announcements/:id/archive | ✅ msg 已实现 | admin-portal 公告归档 |
| REST PUT /announcements/:id/pin | ✅ msg 已实现 | admin-portal 公告置顶 |
需 teacher-bff(ai03):
- 状态更正:ListNotifications / MarkAsRead 已实现
- 公告改用 REST 聚合(不期望 gRPC ListAnnouncements)
- admin-portal 的 4 个公告 mutation(createAnnouncement/publishAnnouncement/archiveAnnouncement/toggleAnnouncementPin)通过 teacher-bff HTTP 调用 msg REST
4.4 student-bff(ai04)— ✅ 全部就绪
student-bff nextstep-v2.md 确认 msg "✅ 13 RPC + REST 就绪"。
4.5 parent-bff(ai05)— ⚠️ 需状态更正 + 命名对齐
| 依赖 | 状态 | 说明 |
|---|---|---|
| gRPC listNotifications | ✅ proto 已定义 | parent-bff 标注 ✅ 正确 |
| gRPC markAsRead | ✅ proto 已定义 | parent-bff 标注 ✅ 正确 |
| gRPC getNotificationPreferences | ✅ proto 已定义 | parent-bff 标注"⚠️ proto 未定义"有误 |
| gRPC updateNotificationPreferences | ✅ proto 已定义 | parent-bff 标注"⚠️ proto 未定义"有误 |
| /healthz | ✅ 已运行 | parent-bff 标注"⏳ 未运行"有误 |
需 parent-bff(ai05):
- 状态更正:NotificationPreferenceService 已定义
- RPC 命名对齐:proto 为
GetPreferences/UpdatePreferences,非getNotificationPreferences/updateNotificationPreferences - /healthz 状态更正为 ✅
4.6 前端 portal(ai13-ai16)— ⚠️ 命名不一致
parent-portal nextstep-v2.md §1 标注命名不一致:
| 前端 operation | 后端命名 | 说明 |
|---|---|---|
myNotifications |
notifications |
前端加 my 前缀 |
myNotificationPreferences |
notificationPreferences |
前端加 my 前缀 |
markAsRead |
markNotificationRead |
命名不一致 |
处理方案:命名不一致由各 BFF 在 resolver 层映射,msg proto/REST 命名保持不变。需各 BFF 确认 resolver 已做字段映射。
5. msg 需要上下游实现的工作
5.1 需要上游(同层级)实现
| 上游模块 | 需求 | 状态 | 说明 |
|---|---|---|---|
| core-edu | 确认 3 个考试实时事件 payload 包含 studentIds |
⏳ 待确认 | msg handler 依赖 studentIds 确定通知接收人;若无则跳过 |
| core-edu | 确认 topic 命名对齐 | ⚠️ 待确认 | msg 订阅 edu.teaching.assignment.submitted/graded,core-edu §4.2 列出 homework.submitted/graded,需确认 topic 名是否一致 |
| iam | Kafka 事件 payload 字段对齐 | ✅ 已确认 | msg 按 camelCase + snake_case 双兼容消费 |
| data-ana | edu.insight.mastery.updated 事件发布 |
⏳ 待 data-ana 就绪 | msg 已订阅,等待 data-ana 发布事件 |
5.2 需要下游实现的工作
| 下游模块 | 需求 | 状态 | 说明 |
|---|---|---|---|
| api-gateway | 新增 /api/v1/announcements/* 路由 → msg:3007 |
⏳ 待 ai01 实现 | 公告 REST API 需通过 gateway 暴露 |
| teacher-bff | 状态更正:ListNotifications/MarkAsRead 已实现 | ⏳ 待 ai03 更新 | nextstep-v2.md §2.5 标注"⏳"有误 |
| teacher-bff | 公告改用 REST 聚合(非 gRPC) | ⏳ 待 ai03 实现 | msg 公告为 REST only,teacher-bff 通过 HTTP 调用 msg:3007/announcements/* |
| parent-bff | 状态更正:NotificationPreferenceService 已定义 | ⏳ 待 ai05 更新 | nextstep-v2.md §4.4 标注"⚠️ proto 未定义"有误 |
| parent-bff | RPC 命名对齐:GetPreferences/UpdatePreferences | ⏳ 待 ai05 更新 | proto 标准命名,非 getNotificationPreferences |
| push-gateway | 消费 edu.notify.notification.sent |
✅ 已对齐 | push-gateway v2 已完成 |
| 各 BFF | GraphQL 字段命名映射 | ⏳ 待各 BFF 确认 | myNotifications vs notifications 等,由 BFF resolver 层处理 |
6. 待协调事项
| # | 事项 | 协调对象 | 说明 |
|---|---|---|---|
| 1 | core-edu 考试实时事件 payload 确认 | ai07 | 确认 edu.teaching.exam.extended/force_submitted/question_reordered 的 payload 是否包含 studentIds: string[]。msg handler 依赖此字段 |
| 2 | core-edu topic 命名对齐 | ai07 | msg 订阅 edu.teaching.assignment.submitted/graded,core-edu §4.2 列出 homework.submitted/graded。需确认实际 topic 名 |
| 3 | api-gateway 新增 announcements 路由 | ai01 | /api/v1/announcements/* → msg:3007 |
| 4 | teacher-bff 状态更正 + 公告 REST 接入 | ai03 | ListNotifications/MarkAsRead 已实现;公告改用 REST |
| 5 | parent-bff 状态更正 + RPC 命名对齐 | ai05 | NotificationPreferenceService 已定义;命名用 GetPreferences/UpdatePreferences |
| 6 | 公告 gRPC RPC 决策 | coord | 下游期望 ListAnnouncements gRPC,但 ARB-008 限制 RPC 总数 13。msg 保持 REST only,需 coord 仲裁 |
| 7 | 考试实时事件责任方澄清 | ai07/ai09 | student-portal 期望 ExamExtended/ExamForceSubmitted/ExamQuestionReordered 事件。core-edu 已发布 3 topic,msg 已消费并转发到 edu.notify.notification.sent,push-gateway 推送到 WebSocket。链路已通 |
7. Docker 部署配置
7.1 启动命令
docker run -d --name edu-msg --network edu-full_default \
-p 3007:3007 -p 50056:50056 \
-e PORT=3007 -e GRPC_PORT=50056 \
-e DATABASE_URL=mysql://edu:changeme@edu-mysql:3306/next_edu_cloud \
-e REDIS_URL=redis://edu-redis:6379 \
-e KAFKA_BROKERS=edu-kafka:29092 \
-e ES_URL=http://edu-es:9200 \
-e PUSH_GATEWAY_URL=http://push-gateway:8081 \
-e DEV_MODE=false -e NODE_ENV=production -e LOG_LEVEL=info \
edu-msg:test
7.2 就绪信号
| 信号 | 状态 | 说明 |
|---|---|---|
| HTTP :3007 /healthz | ✅ | {"status":"ok","service":"msg"} |
| HTTP :3007 /readyz | ✅ | 5/6 OK(pushGateway 软失败,预期) |
| gRPC :50056 13 RPC | ✅ | 3 Service 全部注册 |
| Kafka 消费 16 topic | ✅ | consumer group 加入成功 |
| Kafka 发布 4 topic | ✅ | Outbox 轮询 5s,ARB-013 命名 |
| Docker 镜像 edu-msg:test | ✅ | 多阶段构建,pnpm workspace 模式 |
8. v2 工作完成总结
8.1 已完成的全部 v2 工作
- ARB-013 topic 命名统一:
edu.notification.*→edu.notify.notification.*(4 topic + FALLBACK) - homework.assigned 消费:新增 topic 订阅 + handler
- 3 个考试实时事件消费:exam.extended/force_submitted/question_reordered,转发到 push-gateway
- 嵌套 payload 解包:兼容 core-edu Outbox
{ event_id, payload: {...} }格式 - groupId 长度安全:优先用业务 eventId,降级路径 truncate 到 32 字符
- sendBatch eventId 幂等:批量查询过滤已存在 eventId(v1 完成,v2 Docker 验证)
- pino ESM 导入修复:
import { pino } from 'pino' - Dockerfile pnpm prune 修复:
pnpm install --prod --ignore-scripts - Docker 本地真实环境验证:16 topic 订阅 + 5 条考试实时事件通知入库 + Outbox 转发
- 101 单元测试通过:typecheck + lint 零错误
8.2 msg v2 对外能力
- gRPC:13 RPC(3 Service)
- REST:27 endpoints(notifications 10 + preferences 2 + templates 6 + announcements 9)
- Kafka 消费:16 类事件(iam 6 + core-edu 9 + data-ana 1)
- Kafka 发布:4 类事件(ARB-013 命名)
- Docker:edu-msg:test 镜像就绪,/healthz + /readyz + /metrics 全部可用
msg v2 完成。13 gRPC RPC + 27 REST endpoints + 16 类 Kafka 消费 + 4 类 Kafka 发布全部就绪,101 单元测试通过,Docker 本地真实环境验证全部通过(无 mock 数据)。等待 core-edu 确认考试实时事件 payload 字段 + api-gateway 新增 announcements 路由 + teacher-bff/parent-bff 状态更正后即可端到端联调。