# Edu 平台架构 v2 重设计 > 版本:v2.1 > 日期:2026-07-14 > 状态:待评审 > 变更摘要(v2.1):Temporal 严格边界引入(AI 工作流 + Saga);CDC+Outbox 结合(Debezium Transaction Log Tailing,去掉轮询);DataLoader 强制;Router-Authorization 信任凭证;外部 GraphQL + 内部 gRPC;Eager Invalidation + 乐观锁版本号;Redis Pub/Sub 推送背板;ScopeToken 优化大规模 ID 列表。 > 关联: > > - [004 架构影响地图](../../architecture/004_architecture_impact_map.md) v2.0(当前态) > - [Portal Shell 插件化仪表盘设计](./2026-07-14-portal-shell-widget-dashboard-design.md) v2.1 > - [0010 架构蓝图](../../architecture/0010_architecture.md) > - [项目规则](../../.trae/rules/project_rules.md) --- ## 1. 背景与目标 ### 1.1 v1 架构痛点诊断 当前 Edu 微服务架构(v1/v2.0 of 004)存在以下核心痛点: 1. **BFF 样板代码灾难**:3 个手写 BFF(teacher-bff / student-bff / parent-bff)重复实现 GraphQL Resolver + gRPC 调用,每个字段改动需改 3 处(proto → schema → resolver),×3 灾难 2. **DataScope 跨库 JOIN 悖论**:core-edu 要过滤"管理员可见班级",但用户-学校关系在 iam 库,微服务禁止跨库 JOIN,运行时传 ID 列表性能崩溃 3. **CDC 与 Outbox 文档矛盾**:§1.1a 画 CDC(MySQL→Debezium→Kafka),§7.1 画 Outbox(业务表→Relay Worker→Kafka),看起来像两条路径做同一件事 4. **WebSocket 性能开销**:push-gateway 同时支持 WS/SSE,但 K12 场景以单向通知为主,WS 心跳/握手开销过大 5. **Temporal 滥用边界**:v1 规划引入 Temporal 但未明确边界,短事务若走 Temporal 会被入库开销拖垮吞吐;v2 严格立规矩——仅限 AI 耗时工作流 + Saga 分布式事务补偿 6. **iam 职责过载**:iam 承载认证 + RBAC + 审计 + JWKS + 插件配置(6 张表)+ DataScope,成为"上帝服务" 7. **content 多存储耦合**:content 同时写 MySQL + Neo4j + ES,事务一致性、同步延迟、故障域都是问题 8. **ai 有状态化**:ai 服务内嵌 WorkflowStateStore,边界不清,难以水平扩展 9. **配置散落**:配置散落在 env / DB / code,无统一配置中心 10. **前端 4 portal 重复**:4 个独立 Next.js portal 重复实现 Dockerfile / i18n / auth / layout(已在 portal-shell v2.1 设计稿中解决) ### 1.2 v2 设计目标 1. **BFF 联邦化**:采用 Apollo Federation,消灭手写 BFF 样板代码 2. **DataScope 运行时解析**:通过 @requires 指令在子图间传递可见范围,解决跨库悖论 3. **CDC/Outbox 职责分工**:Outbox 负责领域事件,CDC 负责读模型投影 4. **SSE 优先**:推送默认走 SSE,WS 仅限强双向场景 5. **Temporal 严格边界**:仅限 AI 耗时工作流(生成大纲→知识点→题目→组装试卷)+ 跨服务 Saga 分布式事务补偿(如购买插件:扣减积分服务 + 授权服务);CRUD 短事务绝对禁止 6. **服务边界清晰**:iam 拆分、content CQRS、ai 无状态化 7. **配置中心化**:config-service 统一管理动态配置 8. **单 portal-shell**:前端单容器部署(引用 portal-shell v2.1 设计稿) ### 1.3 约束条件 - **多 AI 协作模式不变**:coord + dev + sre,模块单一负责制 - **技术栈全开放**:可换 DB / MQ / 框架 / 语言 - **数据可清零**:无迁移包袱,允许重建测试数据 - **微服务架构保留**:不为当前规模妥协,为未来扩展性预留 - **AI 分支权限**:v2.1 实施期间,人类决策者已授予 AI 提交其他模块代码与切换分支的权限(仅限 v2 实施期,结束后回收) ### 1.4 非目标 - 不重写已有业务逻辑(仅调整架构边界) - 不替换 MySQL / Redis / Kafka / ClickHouse / Neo4j / Elasticsearch - 不引入 etcd / Consul(K12 场景过重) - 不将 Temporal 用于 CRUD 短事务(作业提交、成绩录入等毫秒~秒级事务走同步调用 + 分布式锁) --- ## 2. 整体架构 ### 2.1 6 层分层架构 ``` L1 用户层 └─ 教师 / 学生 / 家长 / 管理员 L2 微前端层(单容器) └─ portal-shell(Modular Monolith + Micro-kernel,Next.js 15 App Router) · RSC 服务端预取 Config + initialData · dynamic import 懒加载插件 · URL Params + Zustand 状态共享 · SWR 静默刷新配置 · 详见 portal-shell v2.1 设计稿 L3 边缘网关层(Go 1.25 / Gin) ├─ api-gateway(JWT 校验 + 限流 + 熔断 + CORS) └─ realtime-gateway(SSE 优先 + WS 仅监考 + Redis Pub/Sub) L4 GraphQL 联邦层(Apollo Router) └─ apollo-router(自动查询计划 + @requires DataScope 传递) · 替代 3 个 BFF 的聚合职责 · 子图自主暴露 GraphQL L5 业务服务层(NestJS / FastAPI,按 DDD 限界上下文) ├─ iam(认证 + RBAC + JWKS + DataScope)— 拆分后 ├─ config-service(插件配置 + 布局 + 用户偏好)— 新拆分 ├─ core-edu(教学核心 + 教学组织) ├─ content(内容资源,CQRS:MySQL 写 + Neo4j/ES 读投影) ├─ msg(消息通知) ├─ data-ana(数据分析,有状态) └─ ai(LLM 推理,无状态) L6 数据层 ├─ MySQL 8.0(写模型,每服务独占 schema) ├─ Redis 7(缓存 + 会话 + Pub/Sub + 限流 + ai workflow 状态) ├─ ClickHouse 24.3(读模型宽表) ├─ Neo4j 5.20(知识图谱读模型) ├─ Elasticsearch 8.13(题库检索读模型) └─ 配置中心(config-service + Redis 缓存) ``` ### 2.2 服务清单(v2) | 类别 | 服务名 | 语言/框架 | HTTP | gRPC | GraphQL | 限界上下文 | 变更说明 | | ------ | ---------------- | -------------------- | ---- | ----- | ------- | ------------------------ | -------------------------------------------------------- | | 边缘 | api-gateway | Go 1.25 / Gin | 8080 | — | — | 网关 | 保留,移除 BFF 代理路由 | | 边缘 | realtime-gateway | Go 1.25 / Gin | 8081 | — | — | 推送 | 改名,SSE 优先 | | 联邦 | apollo-router | Rust (Apollo Router) | 4011 | — | ✅ | GraphQL 联邦 | 新增(避开 portal 4000-4009 段位,与 v1 并行运行不冲突) | | 业务 | iam | NestJS | 3002 | 50052 | ✅ 子图 | 认证+RBAC+JWKS+DataScope | 拆分,移除插件配置 | | 业务 | config-service | NestJS | 3011 | 50059 | ✅ 子图 | 插件+布局+偏好 | 新增,从 iam 拆出 | | 业务 | core-edu | NestJS | 3004 | 50053 | ✅ 子图 | 教学核心+组织 | 保留 | | 业务 | content | NestJS | 3005 | 50054 | ✅ 子图 | 内容资源 | CQRS 改造 | | 业务 | msg | NestJS | 3007 | 50056 | ✅ 子图 | 沟通通知 | 保留 | | 业务 | data-ana | FastAPI | 3006 | 50055 | ✅ 子图 | 数据分析 | 保留 | | 业务 | ai | FastAPI | 3008 | 50058 | ✅ 子图 | LLM 推理 | 无状态化 | | 微前端 | portal-shell | Next.js 15 | 4010 | — | — | 教师端 Shell | 新增,替代 4 portal | **移除的服务**: - teacher-bff / student-bff / parent-bff → 由 apollo-router 替代 - teacher-portal / student-portal / parent-portal / admin-portal → 由 portal-shell 替代 **移除的组件**: - Debezium Connect → 保留并启用,与 Outbox 结合使用(见 §5.4 / §6.1):业务代码只写业务表 + Outbox 表,去掉 OutboxPublisher 轮询线程,改由 Debezium 监听 binlog 自动将 Outbox 表变更推送到 Kafka(Transaction Log Tailing) - Temporal → 引入但严格立规矩(见 §8.2):仅限 AI 耗时工作流 + Saga 分布式事务补偿,CRUD 短事务绝对禁止 ### 2.3 核心数据流 ```mermaid graph TB User[用户] --> PortalShell[portal-shell :4010
RSC 预取] PortalShell --> ApiGateway[api-gateway :8080
JWT + 限流] ApiGateway --> ApolloRouter[apollo-router :4011
GraphQL 联邦] ApolloRouter -->|子图查询| Iam[iam 子图] ApolloRouter -->|子图查询| Config[config-service 子图] ApolloRouter -->|子图查询| CoreEdu[core-edu 子图] ApolloRouter -->|子图查询| Content[content 子图] ApolloRouter -->|子图查询| Msg[msg 子图] ApolloRouter -->|子图查询| DataAna[data-ana 子图] ApolloRouter -->|子图查询| Ai[ai 子图] ApolloRouter -.@requires DataScope.-> Iam CoreEdu -.接收 visibleClassIds.-> CoreEdu PortalShell -.SSE 推送.-> RealtimeGw[realtime-gateway :8081] RealtimeGw -.Redis Pub/Sub 订阅.-> Msg Msg -.Kafka 消费.-> Kafka Iam <-->|Outbox 表 + Debezium 监听| Kafka[(Kafka)] CoreEdu <-->|Outbox 表 + Debezium 监听| Kafka Content <-->|Outbox 表 + Debezium 监听| Kafka Msg <-->|Outbox 表 + Debezium 监听| Kafka DataAna <-->|Kafka 消费| Kafka Ai <-->|Outbox 表 + Debezium 监听| Kafka Iam --> MySQL[(MySQL)] Config --> MySQL CoreEdu --> MySQL Content --> MySQL Content -.Kafka 投影器.-> Neo4j[(Neo4j)] Content -.Kafka 投影器.-> ES[(Elasticsearch)] Msg --> MySQL DataAna --> ClickHouse[(ClickHouse)] ``` ### 2.4 关键架构决策摘要 | 决策点 | v1 方案 | v2 方案 | 理由 | | ------------ | ------------------- | ------------------------------------------------- | -------------------------------------------------- | | 前端 | 4 portal + MF | 单 portal-shell | Modular Monolith,减少重复 | | BFF | 3 个手写 BFF | Apollo Federation | 工业标准,零样板代码 | | DataScope | 网关透传 + 服务注入 | @requires 运行时解析 | 解决跨库悖论 | | Push | WS + SSE | SSE 优先 + WS 仅监考 | K12 场景单向为主 | | Temporal | 规划引入但无边界 | 严格边界引入(AI 工作流 + Saga),CRUD 短事务禁用 | 长工作流自动重试/休眠/状态持久化,短事务走同步调用 | | iam 职责 | 认证+RBAC+插件配置 | 拆分 iam + config-service | 避免上帝服务 | | content 存储 | MySQL+Neo4j+ES 混写 | MySQL 写 + Outbox + CDC 投影器读模型 | CQRS 分离 | | ai 状态 | 有状态工作流 | 无状态推理 + WorkflowStateStore | 明确边界 | | 配置 | env + DB 散落 | config-service 统一 | 配置中心化 | --- ## 3. Apollo Federation BFF 层 ### 3.1 设计原理 **痛点**:v1 的 3 个手写 BFF 重复实现 GraphQL Resolver + gRPC 调用,每个字段改动需改 3 处(proto → schema → resolver),3 个 BFF = ×3 灾难。 **Apollo Federation 工业标准方案**: - 每个业务服务自主暴露 GraphQL 子图(Subgraph) - Apollo Router 作为唯一入口,自动拆分查询计划(Query Plan) - 客户端发一个 Query,Router 自动决定调用哪些子图、以什么顺序、如何聚合 - 零手写聚合代码 ### 3.2 子图暴露方式:proto → GraphQL 自动生成 从 proto 自动生成 GraphQL schema,业务服务只需实现 Resolver。 **proto 扩展**: ```protobuf // core_edu.proto service ExamService { rpc GetExam(GetExamRequest) returns (Exam) { option (graphql.query) = "exam"; } rpc ListExams(ListExamsRequest) returns (ListExamsResponse) { option (graphql.query) = "exams"; } } message Exam { string exam_id = 1; string title = 2; repeated string class_ids = 3; } ``` **生成的 GraphQL 子图**(core-edu 子图): ```graphql type Query { exam(examId: ID!): Exam exams(classIds: [ID!]!): [Exam!]! } type Exam @key(fields: "examId") { examId: ID! title: String classIds: [String!]! } ``` ### 3.3 各服务子图清单 | 服务 | 子图关键类型 | DataScope 依赖 | | -------------- | ----------------------------------------------------------- | --------------------------------------- | | iam | User / Role / Permission / School | 无(权限源) | | config-service | PluginConfig / LayoutTemplate / UserOverride | 无 | | core-edu | Exam / Homework / Grade / Attendance / Class / Schedule | `visibleClassIds` / `visibleStudentIds` | | content | Textbook / Chapter / Question / KnowledgePoint / LessonPlan | `editableSubjectIds` | | msg | Notification / Template / Preference | `visibleNotificationScopes` | | data-ana | Analytics / Dashboard / Mastery / Warning | `visibleClassIds` / `visibleStudentIds` | | ai | Chat / Question / Expression / LessonPlan / Report | `userId`(个人级) | ### 3.4 Apollo Router 部署架构 ``` ┌─────────────────────────────────────────────────┐ │ api-gateway :8080(Go / Gin) │ │ · JWT RS256 校验 │ │ · 限流 / 熔断 / CORS │ │ · 注入 x-user-id / x-user-role / x-dataScope │ └────────────────┬────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────┐ │ apollo-router :4011(Rust,官方二进制) │ │ · 接收 GraphQL Query │ │ · 解析查询计划(Query Plan) │ │ · 并发调用子图,@requires 传递 DataScope │ │ · 聚合响应返回客户端 │ │ · 内置缓存 / 持续查询 / 遥测 │ └────────────────┬────────────────────────────────┘ │ ┌─────────┼─────────┬─────────┬─────────┐ ▼ ▼ ▼ ▼ ▼ iam:3002 core-edu content msg data-ana /graphql :3004 :3005 :3007 :3006 /graphql /graphql /graphql /graphql ``` **Router 配置**(`router.yaml`): ```yaml supergraph: listen: 0.0.0.0:4011 path: /graphql introspection: true homepage: enabled: true health_check: listen: 0.0.0.0:8088 override_subgraph_url: iam: http://iam:3002/graphql config-service: http://config-service:3011/graphql core-edu: http://core-edu:3004/graphql content: http://content:3005/graphql msg: http://msg:3007/graphql data-ana: http://data-ana:3006/graphql ai: http://ai:3008/graphql headers: all: request: - propagate: named: "authorization" - propagate: named: "x-user-id" - propagate: named: "x-user-role" - propagate: named: "x-dataScope" - propagate: named: "x-request-id" ``` ### 3.5 BFF 样板代码消灭对比 **v1 流程(新增字段 `exam.duration`,手动改动点 6 处)**: 1. 改 `core_edu.proto`,加 `int32 duration = 5;`(手动) 2. 重新生成 proto 代码(自动) 3. 改 teacher-bff 的 `schema.graphql` + resolver(手动) 4. 改 student-bff 的 `schema.graphql` + resolver(手动) 5. 改 parent-bff 的 `schema.graphql` + resolver(手动) 6. 改 portal 的 GraphQL 查询(手动) **v2 流程(手动改动点 2 处)**: 1. 改 `core_edu.proto`,加 `int32 duration = 5;`(手动) 2. 运行 `buf generate`,自动生成 GraphQL schema(自动) 3. core-edu 服务自动暴露 `duration` 字段(自动) 4. Apollo Router 自动拉取新 schema(自动) 5. 改 portal-shell 的 GraphQL 查询(手动) **手动改动点从 6 处 → 2 处**(仅改 proto + 改前端查询)。 ### 3.6 DataLoader 强制(@key Reference Resolver) **强制约束**:所有子图的 `@key` 解析器(Reference Resolver)必须使用 DataLoader 进行请求合并(Batching),禁止 N+1 查询。 **原理**:Apollo Federation 的 `@key` 解析器会被 Router 并发调用(如查 50 个 Exam 各自的 Class 信息,会触发 50 次 `__resolveReference`)。若不用 DataLoader,就是 50 次 DB 查询;用 DataLoader 则合并为 1 次 `WHERE id IN (...)` 查询。 **示例**(core-edu 子图的 Exam @key 解析器): ```typescript import DataLoader from "dataloader"; // 批量查询函数:接收 examId 数组,返回 Exam 数组 async function batchLoadExams( examIds: readonly string[], examRepo: ExamRepository, ): Promise { const exams = await examRepo.findByIds([...examIds]); // 保证顺序与 examIds 一致 return examIds.map((id) => exams.find((e) => e.examId === id)); } @Resolver(() => Exam) export class ExamResolver { private examLoader: DataLoader; constructor( private examRepo: ExamRepository, @Inject("REQUEST_SCOPE") requestScope: RequestScope, ) { // 每个请求独立 DataLoader 实例(避免跨请求缓存污染) this.examLoader = requestScope.getOrInit( "examLoader", () => new DataLoader((ids) => batchLoadExams(ids, this.examRepo)), ); } @ResolveReference() resolveReference(ref: { examId: string }): Promise { return this.examLoader.load(ref.examId); } } ``` **强制校验**: - ESLint 自定义规则:`@ResolveReference` 装饰器内禁止直接调用 Repository,必须走 DataLoader - 单元测试:并发调用 100 次 `resolveReference`,断言 Repository 只被调用 1 次 ### 3.7 Router-Authorization 信任凭证 **强制约束**:Apollo Router 请求子图时必须携带内部信任凭证 `Router-Authorization` Header,各服务的 NestJS Guard 必须拦截并校验该 Header,拒绝任何非 Router 发起的 GraphQL 请求。 **理由**:GraphQL 子图暴露后,若不校验来源,恶意请求可直接绕过 Router 的查询计划与权限聚合,直接调子图。 **实现**: ```yaml # router.yaml — Router 侧自动注入 headers: all: request: - insert: name: "Router-Authorization" value: "${ROUTER_AUTH_SECRET}" # 从环境变量读取 - propagate: named: "authorization" - propagate: named: "x-user-id" - propagate: named: "x-user-role" - propagate: named: "x-dataScope" ``` ```typescript // 子图侧 Guard(各业务服务通用) @Injectable() export class RouterAuthGuard implements CanActivate { constructor(private config: ConfigService) {} canActivate(ctx: ExecutionContext): boolean { const req = ctx.switchToHttp().getRequest(); const routerAuth = req.headers["router-authorization"]; const expected = this.config.get("ROUTER_AUTH_SECRET"); if (routerAuth !== expected) { throw new ForbiddenException("Direct GraphQL access denied; must go through Apollo Router"); } return true; } } // 在 GraphQLModule 中全局应用 @Module({ imports: [ GraphQLModule.forRoot({ autoSchemaFile: true, path: "/graphql", // 仅 /graphql 路径应用 RouterAuthGuard useGuards: [RouterAuthGuard], }), ], }) ``` **配置管理**: - `ROUTER_AUTH_SECRET` 通过 Docker Secret 或 K8s Secret 注入,禁止硬编码 - 开发模式(DEV_MODE=true)可放宽校验,便于本地调试 ### 3.8 外部 GraphQL + 内部 gRPC 边界 **强制约束**: | 调用方向 | 协议 | 用途 | | ----------- | -------------------------------- | -------------------- | | 前端 → 后端 | GraphQL(必须走 Apollo Router) | 外部查询,聚合展现 | | 后端 → 后端 | gRPC(必须直连,禁止走 GraphQL) | 内部 RPC,服务间协作 | **理由**:GraphQL 子图纯粹是为了服务前端展现,不能作为微服务间的内部 RPC 协议。后端互调走 gRPC 的优势: - 性能:gRPC 基于 HTTP/2 + Protobuf,延迟远低于 GraphQL over HTTP/1.1 - 强类型:proto 契约由 buf 强制校验 - 流式:gRPC 原生支持 streaming,GraphQL 不支持 - 不污染 GraphQL schema:内部协作的字段(如跨服务 JOIN 辅助字段)不应暴露给前端 **示例**: - ai → data-ana:gRPC `AnalyticsService.GetStudentMastery`(获取学生掌握度,作为 LLM prompt 上下文) - core-edu → iam:gRPC `IamService.GetUserRoles`(查询用户角色,用于权限校验) - content → core-edu:gRPC `ExamService.GetExamsBySubject`(查询某学科考试列表) **禁止**:ai 子图通过 GraphQL 查询 data-ana 子图的数据(这会让 data-ana 的 GraphQL schema 膨胀,且绕过 gRPC 的性能优势)。 --- ## 4. 业务服务层调整 ### 4.1 iam 拆分:iam + config-service **问题**:iam 承载认证 + RBAC + 审计 + JWKS + 插件配置(v2.1 新增 6 张表)+ DataScope 解析,成为"上帝服务"。 **拆分方案**: | 服务 | 职责 | 表 | gRPC Service | | ---------------------- | ------------------------------------- | ----------------------------------------------------------------------------------------------------------------------- | ----------------------------------------------------- | | iam(瘦身) | 认证 + RBAC + JWKS + 审计 + DataScope | users / roles / permissions / refresh_tokens / sessions / totp / audit_logs / user_school_role | IamService / RbacService / JwksService / AuditService | | config-service(新增) | 插件配置 + 布局配置 + 用户偏好 | plugin_registry / role_plugin_mapping / role_layout_default / layout_templates / user_layout_override / plugin_packages | ConfigService | **拆分理由**: - 认证是"高频低变"(每次请求校验 JWT) - 配置是"低频高变"(admin 改配置、用户改偏好) - 两者耦合会导致:配置表锁竞争影响认证、配置 Schema 变更需重启认证服务 **DataScope 归属**: - iam 保留 `user_school_role` 表(用户-学校-角色关系) - iam 子图暴露 `visibleClassIds` / `visibleStudentIds` 查询(DataScope 解析在 iam) - config-service 不涉及 DataScope **端口分配**: - iam: 3002 / gRPC 50052(不变) - config-service: 3011 / gRPC 50059(新增) ### 4.2 content CQRS 改造 **问题**:content 当前同时写 MySQL + Neo4j + ES,三存储事务一致性、同步延迟、故障域都是问题。 **CQRS 改造方案**: ``` 写侧(Command): content 服务 → MySQL(唯一写模型) ↓ Outbox 事件 Kafka(edu.content.* topic) ↓ 读侧投影器(Query): ├── neo4j-projector(消费事件 → 写 Neo4j 知识图谱) ├── es-projector(消费事件 → 写 ES 题库索引) └── content-cache-projector(消费事件 → 失效 Redis 缓存) ``` **改造点**: - content 服务只写 MySQL,不再直接写 Neo4j/ES - 新增 `neo4j-projector` / `es-projector` 作为 content 服务内 worker - content 服务的 GraphQL 子图读 MySQL(强一致)+ 读 Neo4j/ES(最终一致) **GraphQL 子图设计**: ```graphql type Query { textbook(textbookId: ID!): Textbook # 读 MySQL textbooks: [Textbook!]! # 读 MySQL searchQuestions(keyword: String!): [Question!]! # 读 ES knowledgeGraph(subjectId: ID!): KnowledgeGraph # 读 Neo4j } ``` **好处**: - 写路径简单(单库事务) - 读路径可独立扩展(ES/Neo4j 可单独扩容) - 故障隔离(ES 挂了不影响写入) - 数据一致性通过 Outbox + 投影器保证 ### 4.3 ai 服务无状态化 **问题**:ai 当前有 WorkflowStateStore(lesson_plan_workflow),是有状态的。 **v2 定位**: - ai 是无状态 LLM 推理引擎 - 不存业务数据,不存会话历史 - 会话历史存 Redis(短期)或 data-ana(长期分析) - 工作流状态存 Redis(key: `ai:workflow:{workflowId}`,TTL 1 小时) **改造点**: - 移除 ai 内部 WorkflowStateStore 的持久化(改为 Redis) - lesson_plan_workflow 状态存 Redis - ai 用量事件通过 Outbox 发布到 Kafka(`edu.ai.usage.recorded` topic) - data-ana 消费 ai 用量事件做统计 ### 4.4 data-ana 与 ai 边界 | 维度 | data-ana | ai | | ---- | --------------------------------- | --------------------------------------- | | 定位 | 有状态分析引擎 | 无状态推理引擎 | | 存储 | ClickHouse(宽表)+ Redis(缓存) | 无业务存储(仅 Redis 存 workflow 状态) | | 输入 | 业务事件(Outbox 消费) | 用户请求 + data-ana 数据 | | 输出 | 统计结果 + 预警 + 掌握度 | LLM 生成内容(聊天/题目/报告) | | 交互 | ai 调 data-ana 获取分析数据 | data-ana 不调 ai | **调用关系**: - ai → data-ana:gRPC(获取学生掌握度、班级表现等数据,作为 LLM prompt 上下文) - data-ana → ai:无(data-ana 不依赖 ai) - ai 用量事件 → Kafka → data-ana 消费(统计 token 使用量) --- ## 5. 数据层 + DataScope ### 5.1 存储矩阵(v2) | 存储 | 版本 | 用途 | 使用服务 | 变更 | | ------------- | ---- | ----------------------------------------------- | ----------------------------------------------- | -------------------------- | | MySQL | 8.0 | 写模型主库(每服务独占 schema) | iam / config-service / core-edu / content / msg | 新增 config-service schema | | Redis | 7 | 缓存 / 会话 / 限流 / Pub/Sub / ai workflow 状态 | 全部服务 | 新增 ai workflow key | | ClickHouse | 24.3 | 读模型宽表 / 分析聚合 | data-ana | 不变 | | Neo4j | 5.20 | 知识图谱(content 读模型投影) | content(读) | 改为投影器写入 | | Elasticsearch | 8.13 | 题库检索 / 消息检索(读模型投影) | content / msg(读) | content 改为投影器写入 | | 配置中心 | — | 动态配置(插件/布局/特性开关) | config-service | 新增 | ### 5.2 DataScope @requires 实现 #### K12 场景规模分析 | 角色 | 可见范围 ID 量级 | 性能评估 | | -------- | ---------------- | ------------------------- | | 学生 | 1 个班 | 无压力 | | 教师 | 5-10 个班 | IN 查询可扛 | | 教研组长 | 20-50 个班 | IN 查询可扛 | | 管理员 | 100-500 个班 | 需索引优化 | | 超管 | 全校 | 走 DataScope=ALL,不传 ID | **结论**:K12 场景 ID 列表通常 < 500,@requires 运行时解析完全可行。 #### DataScope 6 级(保留 v1) | Level | 含义 | 实现 | | --------- | ------ | --------------------------------------- | | L1 SELF | 仅自己 | `WHERE user_id = ?` | | L2 CLASS | 本班 | `WHERE class_id IN (visibleClassIds)` | | L3 GRADE | 本年级 | `WHERE grade_id IN (visibleGradeIds)` | | L4 SCHOOL | 本校 | `WHERE school_id IN (visibleSchoolIds)` | | L5 REGION | 本区域 | `WHERE region_id IN (visibleRegionIds)` | | L6 ALL | 全部 | 无 WHERE | #### @requires 查询计划示例 **场景**:教研组长查"本组所有班级的考试成绩" **客户端 Query**: ```graphql query GetGradesForMyClasses($termId: ID!) { gradesForCurrentUser(termId: $termId) { examId title classId avgScore } } ``` **Apollo Router 查询计划**(自动生成,使用 ScopeToken 优化): ``` 步骤 1: 调 iam 子图 scopeToken(userId=x-user-id, scopeType="class") → iam 查 Redis 缓存(key: iam:datascope:class:{userId},TTL 5min) → 未命中则查 user_school_role 表 → 将 ID 列表存入 Redis Set(key: scope:usr:{userId}:cls_scope,TTL 5min) → 返回极短的 scopeToken: "usr:123:cls_scope"(不传全量 ID 数组) 步骤 2: 调 core-edu 子图 gradesForCurrentUser( scopeToken: "usr:123:cls_scope", termId: "2024-spring" ) @requires(fields: "scopeToken") → core-edu 从同机房 Redis SMEMBERS scope:usr:123:cls_scope 拿到 ID 数组 → 查 MySQL: WHERE class_id IN (...) AND term_id = ? → 返回 [Grade, Grade, ...] 步骤 3: 聚合返回客户端 ``` **ScopeToken 优化理由**: - 大规模 ID 列表(如管理员可见 500 个班)直接在 GraphQL 联邦中传递全量数组,会导致 Router ↔ 子图之间的 HTTP 请求 payload 膨胀,影响延迟 - ScopeToken 是极短的引用字符串(如 `usr:123:cls_scope`),通过 @requires 传递,子图拿到 token 后从同机房 Redis 极速 MGET/SMEMBERS 获取实际 ID 数组 - Redis 调用延迟 < 1ms,远小于 HTTP payload 传输 500 个 ID 的开销 - ScopeToken 与 Redis Set 共享 TTL(5min),自动过期清理 #### iam 子图 DataScope Resolver(返回 ScopeToken) ```typescript @Resolver(() => User) export class DataScopeResolver { constructor( private iamService: IamService, @Inject("REDIS") private redis: Redis, ) {} @ResolveField(() => String) async classScopeToken( @Parent() user: User, @Context() ctx: { userId: string; dataScope: string }, ): Promise { if (ctx.dataScope === "ALL") return "ALL"; const token = `usr:${ctx.userId}:cls_scope`; const redisKey = `scope:${token}`; // 检查缓存是否存在 const exists = await this.redis.exists(redisKey); if (!exists) { // 未命中则查 DB 并写入 Redis Set const classIds = await this.iamService.getVisibleClassIds(ctx.userId); if (classIds.length > 0) { await this.redis.sadd(redisKey, ...classIds); await this.redis.expire(redisKey, 300); // 5min TTL } } return token; } } ``` #### core-edu 子图 @requires Resolver(从 ScopeToken 解析 ID) ```typescript @Resolver(() => Grade) export class GradeResolver { constructor( private gradeService: GradeService, @Inject("REDIS") private redis: Redis, ) {} @Query(() => [Grade]) @RequirePermission("GRADES_READ") async gradesForCurrentUser( @Context() ctx: { userId: string; dataScope: string; classScopeToken?: string; }, @Args("termId") termId: string, ): Promise { if (ctx.dataScope === "ALL" || ctx.classScopeToken === "ALL") { return this.gradeService.findByTerm(termId); } // 从 Redis Set 批量获取实际 ID const redisKey = `scope:${ctx.classScopeToken}`; const classIds = await this.redis.smembers(redisKey); return this.gradeService.findByClassIds(classIds, termId); } } ``` ### 5.3 配置中心:config-service + Redis **不引入 etcd/Consul**(K12 场景过重),用 config-service + Redis 实现轻量配置中心。 ``` Admin 改配置 ↓ config-service 写 MySQL(plugin_registry 等) ↓ config-service 发 Kafka 事件(edu.config.entry.changed) ↓ 各服务消费事件,失效本地 Redis 缓存 ↓ 下次查询从 config-service 拉新值,回填 Redis ``` **缓存策略**: - Redis key: `config:{type}:{key}`(如 `config:plugin:grades-widget`) - TTL: 5 分钟(兜底失效) - 事件驱动失效(Kafka `edu.config.entry.changed`) **配置类型**: - 插件配置(plugin_registry / role_plugin_mapping) - 布局配置(layout_templates / role_layout_default / user_layout_override) - 特性开关(feature_flags,如 `enable_ai_tutor`) - 系统配置(system_config,如 `max_exam_duration`) ### 5.4 Outbox + CDC 结合(Debezium Transaction Log Tailing) **核心原则**:业务代码依然只写业务表 + Outbox 表(保证业务语义),但去掉 OutboxPublisher 轮询线程,改由 Debezium 监听 binlog 自动将 Outbox 表变更推送到 Kafka(Transaction Log Tailing 模式)。 **架构**: ``` 写侧(业务服务): 业务事务内写业务表 + outbox 表(原子) ↓ 不再有轮询线程 ↓ MySQL binlog(row mode) ↓ Debezium Connect 监听(伪装成 MySQL Slave) ↓ Kafka(edu... topic) ↓ ├── 投影器(CQRS 读模型同步) │ ├── neo4j-projector(content → Neo4j 知识图谱) │ ├── es-projector(content → ES 题库索引) │ ├── content-cache-projector(content → Redis 缓存失效,兜底) │ └── msg-cache-projector(msg → Redis 缓存失效,兜底) │ └── 业务消费者(领域事件处理) ├── msg(消费 edu.teaching.* / edu.identity.* → 发通知) ├── data-ana(消费 edu.teaching.* / edu.ai.* → 写 ClickHouse 宽表) └── realtime-gateway(消费 edu.notify.* → Redis Pub/Sub 推送) ``` **Outbox 表设计**(保留 v1): ```sql CREATE TABLE outbox ( event_id VARCHAR(36) PRIMARY KEY, aggregate_type VARCHAR(50) NOT NULL, aggregate_id VARCHAR(64) NOT NULL, event_type VARCHAR(100) NOT NULL, payload JSON NOT NULL, created_at TIMESTAMP NOT NULL, INDEX idx_aggregate (aggregate_type, aggregate_id), INDEX idx_created (created_at) ); ``` **Debezium Connector 配置**(每个业务服务一个 connector): ```json { "name": "iam-outbox-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "database.hostname": "mysql", "database.port": "3306", "database.user": "debezium", "database.password": "${DEBEZIUM_PASSWORD}", "database.server.id": "184054", "database.include.list": "iam", "table.include.list": "iam.outbox", "database.history.kafka.bootstrap.servers": "kafka:29092", "database.history.kafka.topic": "schema-changes.iam", "transforms": "unwrap", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones": "false", "topic.prefix": "iam", "topic.creation.default.replication.factor": "1" } } ``` **Topic 路由**:Debezium 默认按 `.` 输出,用 SMT 重写为业务语义 topic: ``` iam.outbox → edu.identity.user.created / edu.identity.user.updated / ... core_edu.outbox → edu.teaching.exam.published / edu.teaching.homework.assigned / ... content.outbox → edu.content.question.created / edu.content.knowledge_point.linked / ... ``` **投影器实现位置**:作为 content 服务内的独立 worker(`content/src/workers/`) ``` content/src/ ├── textbooks/ ├── questions/ ├── knowledge-points/ ├── graphql/ # GraphQL 子图(读 MySQL + Neo4j + ES) └── workers/ # 投影器(消费 Kafka → 写读模型) ├── neo4j-projector.worker.ts ├── es-projector.worker.ts └── cache-projector.worker.ts ``` **与 Eager Invalidation 配合**: - 写请求在事务提交后同步发 Redis DEL(主动失效,见 §5.5) - content-cache-projector 作为兜底(防止同步 DEL 因网络抖动丢失) ### 5.5 Eager Invalidation(主动失效) **核心策略**:在写 MySQL 事务提交后的代码行,直接同步发一条 Redis DEL 命令删除查询缓存(而不是等 Kafka 绕一圈回来)。 **双层失效机制**: ``` 写请求 ↓ MySQL 事务提交(业务表 + Outbox 表原子) ↓ 同步行(事务提交后立即执行) Redis DEL cache:content:question:* ← 主动失效(毫秒级) ↓ (异步行,Kafka 绕一圈) content-cache-projector 消费 Kafka → Redis DEL ← 兜底清理(防止网络抖动丢失) ``` **代码示例**: ```typescript @Injectable() export class QuestionService { constructor( private questionRepo: QuestionRepository, @Inject("REDIS") private redis: Redis, private outbox: OutboxService, ) {} async updateQuestion(id: string, dto: UpdateQuestionDto): Promise { return this.questionRepo .transaction(async (trx) => { // 1. 写业务表 const updated = await this.questionRepo.update(id, dto, trx); // 2. 写 Outbox 表(同事务原子) await this.outbox.append( { aggregateType: "Question", aggregateId: id, eventType: "QuestionUpdated", payload: { id, ...dto }, }, trx, ); return updated; }) .then(async (result) => { // 3. 事务提交后,同步发 Redis DEL(Eager Invalidation) // 注意:在事务外执行,避免 Redis 故障导致业务回滚 await this.redis.del(`cache:content:question:${id}`); await this.redis.del(`cache:content:question:list:*`); return result; }); } } ``` **关键原则**: - Redis DEL 在事务**外**执行(事务提交后),避免 Redis 故障导致业务回滚 - 若 Redis DEL 失败,不影响业务结果(兜底由 content-cache-projector 处理) - Kafka 的 content-cache-projector 仅作为防止网络抖动的兜底清理,不是主失效路径 ### 5.6 乐观锁 / 版本号回传 **核心策略**:写接口返回记录的 `updated_at` 或 `version`,前端携带该 version 发起 Query,如果 Router 发现 ES 的数据版本低于预期,强制穿透读一次 MySQL。 **解决的问题**:CQRS 读模型(ES/Neo4j)通过 Outbox + Debezium 异步同步,存在秒级延迟。写后立即读(Read-After-Write)可能读到旧数据。 **流程**: ``` 1. 前端发起写请求(UpdateQuestion) ↓ 2. content 服务写 MySQL,返回 { id, version: 5, updated_at: "2026-07-14T10:00:00Z" } ↓ 3. 前端发起读请求(GetQuestion),携带 If-None-Match: version=5 ↓ 4. Apollo Router 将 version 透传到 content 子图 ↓ 5. content 子图先读 ES(快),检查 ES 中的 version ├── ES.version >= 5 → 返回 ES 数据(最终一致) └── ES.version < 5 → 穿透读 MySQL(强一致),返回 MySQL 数据 ↓ 6. 返回客户端 ``` **GraphQL Schema 扩展**: ```graphql type Question { id: ID! content: String version: Int! updatedAt: DateTime! } input QuestionQueryInput { id: ID! """ 客户端持有的最新版本号,用于读后一致性校验 """ expectedVersion: Int } type Query { question(input: QuestionQueryInput!): Question } ``` **Resolver 实现**: ```typescript @Query(() => Question) async question( @Args("input") input: QuestionQueryInput, ): Promise { // 先读 ES const esResult = await this.esService.get(input.id); if (esResult && (!input.expectedVersion || esResult.version >= input.expectedVersion)) { return esResult; // ES 版本满足,直接返回 } // ES 版本落后或不存在,穿透读 MySQL return this.questionRepo.findById(input.id); } ``` **适用场景**: - 用户编辑后立即刷新页面查看 - 批量导入后立即查询列表 - 不适用于:高频查询场景(version 校验开销),走最终一致即可 ### 5.7 缓存层级与失效策略总结 | 缓存类型 | Redis Key 前缀 | 失效方式 | TTL | | -------------------------- | ----------------- | ------------------------- | ----- | | 查询缓存(content) | `cache:content:*` | Eager DEL + 投影器兜底 | 10min | | 查询缓存(msg) | `cache:msg:*` | Eager DEL + 投影器兜底 | 10min | | DataScope 缓存(iam) | `iam:datascope:*` | 角色变更时 DEL | 5min | | ScopeToken 缓存(iam) | `scope:usr:*` | 自动过期 | 5min | | 配置缓存(config-service) | `config:*` | Kafka 事件失效 + TTL 兜底 | 5min | | AI workflow 状态 | `ai:workflow:*` | 自动过期 | 1h | --- ## 6. 事件驱动架构 ### 6.1 Outbox + CDC 结合模式(Transaction Log Tailing) **核心决策**:v2.1 不再在 Outbox 和 CDC 之间二选一,而是**结合使用**——业务代码只写 Outbox 表(业务语义),由 Debezium 监听 binlog 自动投递到 Kafka(技术传输)。 ``` 写侧:业务表 + Outbox 表(同事务原子) 业务服务(iam/core-edu/content/msg/ai) ① 业务事务内写业务表 + outbox 表(原子) ② 不再有 OutboxPublisher 轮询线程(v1 的轮询模式被废弃) ③ Debezium Connect 监听 MySQL binlog(伪装成 Slave) ④ Debezium 监听到 outbox 表有新增 → 自动推送到 Kafka ⑤ 保证 at-least-once 语义(Debezium offset 管理) 优势:业务代码不感知 MQ;无损监听 binlog 不影响业务性能;消除轮询线程的延迟和 DB 压力 读侧:投影器消费 Kafka(CQRS 读模型同步) ├── neo4j-projector(content → Neo4j 知识图谱) ├── es-projector(content → ES 题库索引) ├── content-cache-projector(content → Redis 缓存失效,兜底) └── msg-cache-projector(msg → Redis 缓存失效,兜底) 用途:CQRS 读模型投影,非业务事件发布 ``` **关键原则**: - 业务代码只写业务表 + Outbox 表(业务语义由业务代码控制) - Outbox 表的投递由 Debezium 自动完成(技术传输由基础设施完成) - 投影器消费 Kafka 事件做 CQRS 读模型同步(与业务事件发布解耦) - Outbox + CDC 不矛盾,而是分工:Outbox 定义"发什么",CDC 解决"怎么发" - v1 的 OutboxPublisher 轮询线程在 v2.1 完全废弃(被 Debezium 替代) ### 6.2 Kafka Topic 命名规范 **格式**:`edu...` | 域 | Topic | 生产者 | 消费者 | 用途 | | -------- | ------------------------------------ | -------------- | ------------------------------ | -------------------- | | identity | `edu.identity.user.created` | iam | msg / core-edu / content | 用户创建 | | identity | `edu.identity.user.updated` | iam | msg / core-edu / content | 用户更新 | | identity | `edu.identity.user.role_changed` | iam | msg / data-ana | 角色变更 | | identity | `edu.identity.role.created` | iam | msg | 角色创建 | | teaching | `edu.teaching.exam.published` | core-edu | msg / data-ana | 考试发布 | | teaching | `edu.teaching.exam.extended` | core-edu | msg | 考试延时 | | teaching | `edu.teaching.homework.assigned` | core-edu | msg / data-ana | 作业布置 | | teaching | `edu.teaching.assignment.submitted` | core-edu | msg / data-ana | 作业提交 | | teaching | `edu.teaching.assignment.graded` | core-edu | msg / data-ana | 作业批改 | | teaching | `edu.teaching.grade.recorded` | core-edu | msg / data-ana | 成绩录入 | | teaching | `edu.teaching.attendance.recorded` | core-edu | msg / data-ana | 考勤记录 | | content | `edu.content.question.created` | content | es-projector / neo4j-projector | 题目创建 | | content | `edu.content.question.updated` | content | es-projector / neo4j-projector | 题目更新 | | content | `edu.content.knowledge_point.linked` | content | neo4j-projector | 知识点关联 | | notify | `edu.notify.notification.sent` | msg | realtime-gateway | 通知发送 | | notify | `edu.notify.notification.read` | msg | — | 通知已读 | | notify | `edu.notify.notification.recalled` | msg | realtime-gateway | 通知撤回 | | insight | `edu.insight.mastery.updated` | data-ana | msg | 掌握度更新 | | ai | `edu.ai.usage.recorded` | ai | data-ana | AI 用量记录 | | config | `edu.config.entry.changed` | config-service | 全部服务 | 配置变更(失效缓存) | **新增 topic**: - `edu.ai.usage.recorded`:ai 用量事件(回流 data-ana 统计) - `edu.config.entry.changed`:配置变更通知(失效各服务本地缓存) ### 6.3 事件版本化 - Schema Registry:用 Kafka 的 schema registry(或简化为 proto 版本号) - 版本后缀:`v1` / `v2`(如 `edu.teaching.exam.published.v2`) - 向后兼容:新增字段用 optional,删除字段用 reserved - 破坏性变更:升版本号,新旧 topic 并行,消费者逐步迁移 ### 6.4 幂等性 | 组件 | 幂等机制 | | ---------------- | ----------------------------------------------------------------------------------------------------- | | Kafka Producer | `idempotent=true` + `transactionalId` | | Outbox Publisher | v2.1 已废弃轮询模式,改由 Debezium 监听 binlog 自动投递(基于 binlog position offset,at-least-once) | | Consumer | 基于 `event_id` 去重(DB 唯一索引或 Redis SETNX) | | 投影器 | 基于 `event_id` + `aggregate_id` 去重 | --- ## 7. 认证与权限 ### 7.1 JWT 流程(v2) ``` 用户登录 ↓ api-gateway → iam /auth/login ↓ iam 校验密码 → 签发 JWT RS256(含 userId / role / dataScope) ↓ api-gateway 设置 httpOnly Cookie(JWT) ↓ portal-shell 后续请求携带 Cookie ↓ api-gateway 校验 JWT → 注入 x-user-id / x-user-role / x-dataScope / x-request-id ↓ apollo-router 透传 x-user-* 头到子图 ↓ 子图 Resolver 从 context 读取 userId / dataScope ↓ iam 子图 DataScope Resolver 查 visibleClassIds(Redis 缓存) ↓ core-edu 子图 @requires 接收 visibleClassIds,注入 WHERE ``` ### 7.2 权限校验分层 | 层级 | 职责 | 实现 | | -------------- | ------------------------------ | ---------------------------- | | api-gateway | JWT 校验 + 限流 | Go 中间件 | | apollo-router | 透传 x-user-* 头 | 配置 headers.propagate | | 子图 Resolver | @RequirePermission 装饰器 | NestJS Guard | | 子图 DataScope | @requires 注入 visibleClassIds | Apollo Federation 指令 | | Repository | WHERE 注入 visibleClassIds | DataScopeInjector(保留 v1) | ### 7.3 DEV_MODE - `DEV_MODE=true` 跳过 JWT 校验,接受 `dev-token` - 保留现有 dev-token 预定义角色机制 --- ## 8. 可观测性 + 工作流 + 推送 ### 8.1 可观测性 | 支柱 | 组件 | v2 变更 | | -------- | -------------------------------------------- | -------------------------- | | 日志 | pino / zap / structlog | 不变 | | 指标 | prom-client / prometheus / prometheus-client | 新增 apollo-router 指标 | | 链路 | OpenTelemetry + Jaeger 1.57 | 不变 | | 健康检查 | /healthz + /readyz | 新增 apollo-router /health | **apollo-router 可观测性**: - 内置 Prometheus 指标(请求量/延迟/错误率/子图延迟) - 内置 OpenTelemetry trace - 内置持续查询(Persisted Queries)统计 ### 8.2 工作流边界(Temporal 严格约束) **v2.1 决策**:引入 Temporal,但严格立规矩——仅限 AI 耗时工作流 + Saga 分布式事务补偿,CRUD 短事务绝对禁止。 **架构约束**: | 工作流类型 | 时长 | 推荐方案 | 示例 | | ------------------- | --------- | -------------------------------------------- | ------------------------------------------------------ | | AI 耗时工作流 | 分钟~小时 | **Temporal Workflow** | 生成大纲 → 生成知识点 → 生成题目 → 组装试卷 | | Saga 分布式事务补偿 | 秒~分钟 | **Temporal Saga** | 购买高级插件:扣减积分服务(+补偿)+ 授权服务(+补偿) | | 长工作流(非 AI) | 分钟~天 | Redis 状态机 + Outbox 事件 | 备课工作流(非 AI 部分)、报告生成 | | 短事务 | 毫秒~秒 | **同步调用 + 分布式锁**(绝对禁止 Temporal) | 作业提交、成绩录入、考勤记录 | | 批处理 | 小时 | Cron + Batch Job | 成绩统计、掌握度计算 | **禁止**: - **CRUD 短事务绝对禁止用 Temporal**:作业提交、成绩录入、考勤记录等毫秒~秒级事务走同步调用 + 分布式锁,Temporal 的入库开销会拖垮吞吐量 - 纯读操作走工作流(直接走缓存) **允许**: - AI 耗时工作流用 Temporal(自动重试 / 休眠 / 状态持久化) - 跨服务 Saga 分布式事务补偿用 Temporal(自动补偿 / 状态可查询) - 跨天长流程用 Redis 状态机 + Outbox 事件(非 AI 场景) **Temporal 部署**: - 独立 Temporal Server 集群(不与业务服务混部) - ai 服务作为 Temporal Worker,注册 Workflow 与 Activity - Workflow 状态持久化在 Temporal DB(PostgreSQL),业务数据仍在各服务 DB ### 8.3 推送架构(SSE 优先 + Redis Pub/Sub 背板) **核心原则**:边缘网关(realtime-gateway)不直接挂载 Kafka,避免在边缘节点引入重量级 Kafka Client。采用 Redis Pub/Sub 作为状态路由背板。 ``` msg 服务发布通知 ↓ Outbox 表 + Debezium → Kafka(edu.notify.notification.*) ↓ msg 服务的后台 worker 消费 Kafka · 持久化通知记录 · 脱敏处理 · 分类聚合 ↓ msg worker 将需要实时推送的消息,按 userId 发到 Redis Pub/Sub · Channel: user:{userId}:notify · Payload: { notificationId, type, title, body, createdAt } ↓ realtime-gateway 实例仅在用户接入 SSE 时,才动态 SUBSCRIBE 该用户的 Redis Channel · 用户上线 → SUBSCRIBE user:{userId}:notify · 用户下线 → UNSUBSCRIBE user:{userId}:notify ↓ ├── SSE 推送(默认,单向) │ · 客户端 GET /sse?token=JWT,保持长连接 │ · realtime-gateway 收到 Redis Pub/Sub 消息 → 推送到 SSE 连接 │ · 适合 99% 场景(考试发布/成绩推送/通知) │ └── WebSocket 推送(仅监考场景) · 客户端 WS /ws,双向通信 · 用于在线监考(心跳/防作弊/实时指令) · MVP 可不实现,二期按需 ``` **架构优势**: - 边缘网关轻量:只依赖 Redis Client,不依赖 Kafka Client(内存占用降低 80%+) - 按需订阅:realtime-gateway 仅订阅在线用户的 Channel,离线用户的消息由 msg worker 持久化待取 - 水平扩展:realtime-gateway 多实例时,Redis Pub/Sub 自动广播到所有订阅实例(用户连哪个实例都能收到) - 故障隔离:Kafka 故障不影响已建立 SSE 连接的实时推送(Redis Pub/Sub 独立) **realtime-gateway 端点**: - `GET /sse` — SSE 推送(默认),连接时 SUBSCRIBE Redis Channel - `GET /ws` — WebSocket 升级(监考专用) - `GET /online/:userId` — 查询在线状态(检查 Redis SET `online:users`) - `POST /internal/broadcast` — 管理员全量广播(特殊场景,如系统维护通知) - `GET /healthz` / `GET /readyz` / `GET /metrics` **在线状态管理**: - 用户 SSE 连接建立 → `SADD online:users {userId}` + 设置 TTL - 用户 SSE 连接断开 → `SREM online:users {userId}` - msg worker 查询在线状态 → `SISMEMBER online:users {userId}`,在线则发 Redis Pub/Sub,离线则仅持久化 **性能对比**: - SSE:单连接内存 ~10KB,可扛 10 万连接/节点 - WS:单连接内存 ~100KB,可扛 1 万连接/节点 - Redis Pub/Sub:单 Channel 推送延迟 < 1ms,百万级 Channel 可承载 --- ## 9. 架构约束(v2.1 新增 + v2 保留) ### 9.1 v2.1 新增约束 | 约束 | 说明 | | ---------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------- | | **Temporal 严格边界** | 仅限 AI 耗时工作流 + Saga 分布式事务补偿;CRUD 短事务(作业提交/成绩录入/考勤)绝对禁止用 Temporal,走同步调用 + 分布式锁 | | **Outbox + CDC 结合** | 业务代码只写业务表 + Outbox 表;废弃 OutboxPublisher 轮询线程;Debezium 监听 binlog 自动投递 Outbox 表到 Kafka(Transaction Log Tailing) | | **DataLoader 强制** | 所有子图 @key 解析器必须使用 DataLoader 请求合并,禁止 N+1 查询;ESLint 规则强制校验 | | **Router-Authorization** | Apollo Router 请求子图必须携带 Router-Authorization Header;各服务 NestJS Guard 校验,拒绝非 Router 的直接 GraphQL 请求 | | **外部 GraphQL + 内部 gRPC** | 前端 → 后端走 GraphQL(必须经 Apollo Router);后端 → 后端走 gRPC(禁止走 GraphQL 子图) | | **Eager Invalidation** | 写事务提交后同步发 Redis DEL(主动失效),Kafka 投影器仅作兜底清理 | | **乐观锁版本号回传** | 写接口返回 version/updated_at;读请求携带 expectedVersion,ES 版本落后则穿透读 MySQL | | **Redis Pub/Sub 推送背板** | 边缘网关不直接挂 Kafka;msg worker 消费 Kafka → 发 Redis Pub/Sub → realtime-gateway 订阅推送 | | **ScopeToken 优化** | 大规模 ID 列表不传全量数组,传 ScopeToken;子图从 Redis SMEMBERS 获取实际 ID | ### 9.2 v2 保留约束 | 约束 | 说明 | | ----------------------- | --------------------------------------------------------------------------------- | | **BFF 联邦化** | 禁止手写 BFF 聚合层,所有 GraphQL 查询走 Apollo Router | | **子图自主** | 每个业务服务必须暴露 GraphQL 子图,proto → GraphQL 自动生成 | | **DataScope @requires** | 禁止跨库 JOIN,DataScope 通过 @requires 指令运行时解析(v2.1 用 ScopeToken 优化) | | **Outbox 领域事件** | 领域事件必须走 Outbox 表,禁止直接调 Kafka producer | | **SSE 优先** | 推送默认走 SSE,WS 仅限强双向场景(监考) | | **iam 职责限定** | iam 仅负责认证/RBAC/JWKS/审计/DataScope,禁止承载配置 | | **config 统一** | 动态配置必须走 config-service,禁止散落 env/code | | **content CQRS** | content 仅写 MySQL,Neo4j/ES 通过投影器同步 | | **ai 无状态** | ai 不存业务数据,会话/工作流状态存 Redis(AI 耗时工作流用 Temporal) | | **单 portal-shell** | 前端单容器部署,禁止新增独立 portal | ### 9.3 保留约束(v1) - 四层分层:Gateway → Router → Services → Data - 依赖方向单向,禁止反向依赖 - 服务间通信:GraphQL(@requires,仅前端)+ gRPC(服务间)+ Kafka(事件) - 每服务独占 DB schema - JWT RS256 + JWKS - Cookie: httpOnly + Secure + SameSite=Strict - buf v2 + FILE 级 breaking 检查 --- ## 10. 迁移路径 ### 10.1 阶段划分(从 v1 到 v2.1) | 阶段 | 内容 | 验收标准 | 依赖 | | ---- | --------------------------------------------------------- | ----------------------------------------- | ------- | | M0 | proto → GraphQL 代码生成工具链 | `buf generate` 输出 GraphQL schema | 无 | | M0.5 | Debezium Connect 部署 + Outbox connector | Debezium 监听 outbox 表推送到 Kafka | 无 | | M1 | 各服务暴露 GraphQL 子图 + DataLoader + RouterAuthGuard | `/graphql` 端点可查询,非 Router 请求被拒 | M0 | | M2 | apollo-router 部署 + Supergraph 组装 | Router 可聚合查询 | M1 | | M3 | iam 拆分 config-service | config-service 独立运行 | M1 | | M4 | DataScope @requires + ScopeToken 实现 | 教师查询可见班级成绩正确 | M2 | | M5 | content CQRS 改造(投影器 + Eager Invalidation + 乐观锁) | Neo4j/ES 通过投影器同步,写后读一致 | M0.5/M1 | | M6 | ai 无状态化 | ai 不存业务数据,状态在 Redis | M1 | | M6.5 | Temporal Server 部署 + ai 接入 Worker | AI 耗时工作流可运行 | M6 | | M7 | realtime-gateway SSE 优先 + Redis Pub/Sub | SSE 推送可用,边缘不挂 Kafka | 无 | | M8 | portal-shell 接入 apollo-router | portal-shell 查询走 Router | M2/M4 | | M9 | 旧 BFF 下线(teacher/student/parent-bff) | 流量为 0 | M8 | | M10 | 旧 portal 下线 | 流量为 0 | M8 | | M11 | 废弃各服务 OutboxPublisher 轮询线程 | 无服务使用轮询投递 | M0.5 | | M12 | arch.db + 004 文档同步 | arch:scan 通过 | 全部 | ### 10.2 并行运行策略 - v1 的 3 个 BFF + 4 个 portal 保留并行运行 - v2 的 apollo-router + portal-shell 新建 - 用户通过路由前缀切换(`/shell/*` 走新,其他走旧) - 逐步迁移用户,流量切完后下线旧服务 ### 10.3 回滚策略 - 任意阶段失败,回滚到上一阶段 - v1 服务始终可用,v2 失败不影响现有用户 - iam 拆分 config-service 时,config 表与现有表无外键依赖,可独立回滚 --- ## 11. ADR 记录 | ADR | 决策 | 理由 | | ------- | ---------------------------------------- | ---------------------------------------------------------------------------------------- | | ADR-023 | 采用 Apollo Federation 替代 3 个手写 BFF | 工业标准,消灭样板代码 | | ADR-024 | DataScope 通过 @requires 运行时解析 | 解决跨库 JOIN 悖论 | | ADR-025 | proto → GraphQL 自动生成 | 单一数据源,字段变更只改 proto | | ADR-026 | iam 拆分 config-service | 避免上帝服务 | | ADR-027 | content CQRS 改造 | 写读分离,故障隔离 | | ADR-028 | ai 无状态化 | 明确边界,可扩展 | | ADR-029 | SSE 优先 + WS 仅监考 | K12 场景单向为主 | | ADR-030 | Temporal 严格边界引入(v2.1 修订) | AI 工作流 + Saga 用 Temporal;CRUD 短事务绝对禁止,走同步调用 + 分布式锁 | | ADR-031 | config-service + Redis 轻量配置中心 | K12 场景,不引入 etcd | | ADR-032 | Outbox + CDC 结合(v2.1 修订) | 业务写 Outbox 表,Debezium 监听 binlog 自动投递(Transaction Log Tailing),废弃轮询线程 | | ADR-033 | portal-shell 单容器替代 4 portal | Modular Monolith | | ADR-034 | realtime-gateway 改名 | SSE 优先定位 | | ADR-035 | DataLoader 强制(@key 解析器) | 消除 N+1 查询,ESLint 规则强制校验 | | ADR-036 | Router-Authorization 信任凭证 | 拒绝非 Router 的直接 GraphQL 请求,防止绕过聚合层 | | ADR-037 | 外部 GraphQL + 内部 gRPC 边界 | GraphQL 仅服务前端展现,后端互调走 gRPC(性能 + 强类型 + 流式) | | ADR-038 | Eager Invalidation(主动失效) | 写后同步 Redis DEL,Kafka 投影器仅作兜底,降低读延迟 | | ADR-039 | 乐观锁版本号回传 | CQRS 读后一致性:写返回 version,读携带 expectedVersion,ES 落后则穿透读 MySQL | | ADR-040 | Redis Pub/Sub 推送背板 | 边缘网关不挂 Kafka,msg worker → Redis Pub/Sub → realtime-gateway 订阅,轻量 + 按需订阅 | | ADR-041 | ScopeToken 优化大规模 ID 列表 | 不传全量数组,传极短 token,子图从 Redis SMEMBERS 获取,降低 HTTP payload 开销 | --- ## 12. 对接清单 ### 12.1 新增服务 | 路径 | 职责 | 动作 | | -------------------------- | -------------------------------------------------- | ---- | | `services/apollo-router/` | Apollo Router 部署配置(router.yaml + Dockerfile) | 新建 | | `services/config-service/` | 配置服务(NestJS,从 iam 拆出) | 新建 | | `services/temporal/` | Temporal Server 部署配置(docker-compose + ui) | 新建 | | `apps/portal-shell/` | Portal Shell 前端(详见 portal-shell v2.1 设计稿) | 新建 | ### 12.2 改造服务 | 路径 | 变更 | 动作 | | ------------------------ | ---------------------------------------------------------------------------------------------- | ---- | | `services/iam/` | 移除插件配置相关表和逻辑到 config-service;新增 GraphQL 子图;废弃 OutboxPublisher 轮询 | 改造 | | `services/core-edu/` | 新增 GraphQL 子图;实现 @requires DataScope(ScopeToken);废弃 OutboxPublisher 轮询 | 改造 | | `services/content/` | CQRS 改造:移除直接写 Neo4j/ES,新增投影器 worker;Eager Invalidation;乐观锁版本号;废弃轮询 | 改造 | | `services/msg/` | 新增 GraphQL 子图;新增 msg worker(消费 Kafka → 发 Redis Pub/Sub);废弃 OutboxPublisher 轮询 | 改造 | | `services/data-ana/` | 新增 GraphQL 子图 | 改造 | | `services/ai/` | 无状态化:WorkflowStateStore 改为 Redis;新增 GraphQL 子图;接入 Temporal Worker(AI 工作流) | 改造 | | `services/push-gateway/` | 改名 realtime-gateway;SSE 优先;移除 Kafka 直连,改 Redis Pub/Sub 订阅 | 改造 | | `services/api-gateway/` | 移除 BFF 代理路由;新增 apollo-router 路由 | 改造 | ### 12.3 下线服务 | 路径 | 处理 | 时机 | | ----------------------- | ----------------------------- | ---- | | `services/teacher-bff/` | 流量切到 apollo-router 后下线 | M9 | | `services/student-bff/` | 流量切到 apollo-router 后下线 | M9 | | `services/parent-bff/` | 流量切到 apollo-router 后下线 | M9 | | `apps/teacher-portal/` | 流量切到 portal-shell 后下线 | M10 | | `apps/student-portal/` | 流量切到 portal-shell 后下线 | M10 | | `apps/parent-portal/` | 流量切到 portal-shell 后下线 | M10 | | `apps/admin-portal/` | 流量切到 portal-shell 后下线 | M10 | ### 12.4 新增工具链 | 路径 | 职责 | 动作 | | -------------------------------------- | ------------------------------------------- | ---- | | `packages/shared-proto/gen-graphql.ts` | proto → GraphQL schema 生成器 | 新建 | | `packages/shared-proto/buf.gen.yaml` | 新增 GraphQL 输出配置 | 修改 | | `packages/shared-ts/src/dataloader/` | 通用 DataLoader 工厂(@key 解析器批量加载) | 新建 | | `packages/shared-ts/src/guards/` | RouterAuthGuard 通用实现 | 新建 | ### 12.5 基础设施 | 路径 | 变更 | 动作 | | --------------------------------- | ------------------------------------------------------------------------------------ | ---- | | `infra/docker-compose.deploy.yml` | 新增 apollo-router / config-service / temporal / portal-shell;移除 3 BFF + 4 portal | 修改 | | `infra/port-allocation.md` | 新增 3011 / 4010 / 4011 / 7233(temporal)端口 | 修改 | | `infra/init-sql/` | 新增 config-service schema DDL;新增 temporal DB DDL | 新增 | | `infra/debezium/` | 新增 Debezium Connect 配置(每业务服务一个 outbox connector,SMT 重写 topic) | 新建 | | `infra/docker-compose.tools.yml` | 新增 temporalio/server + temporalio/ui 镜像预拉 | 修改 | --- ## 13. 验收标准 ### 13.1 功能验收 - [ ] Apollo Router 可聚合 7 个子图查询 - [ ] DataScope @requires + ScopeToken 正确传递可见范围 - [ ] iam 拆分后 config-service 独立运行 - [ ] content CQRS 投影器正确同步 Neo4j/ES - [ ] Eager Invalidation:写后同步 Redis DEL 生效 - [ ] 乐观锁版本号:写后读一致(ES 落后则穿透读 MySQL) - [ ] ai 无状态化,工作流状态在 Redis - [ ] Temporal AI 工作流可运行(生成大纲→知识点→题目→组装试卷) - [ ] CRUD 短事务(作业提交/成绩录入)不走 Temporal,走同步调用 - [ ] realtime-gateway SSE 推送可用,边缘不挂 Kafka - [ ] Redis Pub/Sub 推送背板:msg worker → Redis Pub/Sub → realtime-gateway - [ ] Debezium 监听 outbox 表自动投递 Kafka(无轮询线程) - [ ] 所有 @key 解析器使用 DataLoader(无 N+1 查询) - [ ] RouterAuthGuard 拒绝非 Router 的直接 GraphQL 请求 - [ ] 后端互调走 gRPC,不走 GraphQL 子图 - [ ] portal-shell 接入 apollo-router 查询正常 - [ ] 配置变更通过 Kafka 事件失效各服务缓存 ### 13.2 非功能验收 - [ ] `pnpm run lint` + `pnpm run typecheck` 零错误 - [ ] Apollo Router 查询延迟 P99 < 100ms(不含子图) - [ ] DataScope @requires + ScopeToken 查询延迟 P99 < 500ms(含子图) - [ ] DataLoader 批量加载:并发 100 次 @key 解析,Repository 调用 ≤ 1 次 - [ ] SSE 单节点可扛 10 万连接 - [ ] Redis Pub/Sub 推送延迟 < 1ms - [ ] Debezium 投递延迟 < 1s(binlog → Kafka) - [ ] arch.db 更新,004 文档同步 - [ ] 所有 ADR 记录完整 --- ## 14. 参考资料 - [Apollo Federation 官方文档](https://www.apollographql.com/docs/federation/) - [Apollo Router 配置](https://www.apollographql.com/docs/router/configuration) - [GraphQL @requires 指令](https://www.apollographql.com/docs/federation/federated-types/federated-directives/#requires) - [CQRS 模式](https://martinfowler.com/bliki/CQRS.html) - [Outbox 模式](https://microservices.io/patterns/data/transactional-outbox.html) - [SSE vs WebSocket](https://html.spec.whatwg.org/multipage/server-sent-events.html) - [Portal Shell v2.1 设计稿](./2026-07-14-portal-shell-widget-dashboard-design.md) - [004 架构影响地图](../../architecture/004_architecture_impact_map.md) - [项目规则](../../.trae/rules/project_rules.md) --- ## 15. 版本演进对比 | 维度 | v1(当前) | v2.1(本设计) | | --------------------------------------- | ------------------------------------------------------ | ------------------------------------------------------------------------- | | 前端 | 4 portal + MF | 单 portal-shell | | BFF | 3 个手写 BFF | Apollo Federation + DataLoader + Router-Authorization | | DataScope | 网关透传 + 服务注入 | @requires + ScopeToken 优化(Redis Set 引用) | | Push | WS + SSE | SSE 优先 + WS 仅监考 + Redis Pub/Sub 背板(边缘不挂 Kafka) | | Temporal | 规划引入但无边界 | 严格边界引入(AI 工作流 + Saga),CRUD 短事务禁用 | | Outbox 投递 | OutboxPublisher 轮询线程 | Debezium Transaction Log Tailing(监听 binlog 自动投递,废弃轮询) | | iam | 认证+RBAC+插件配置 | 拆分 iam + config-service | | content | MySQL+Neo4j+ES 混写 | CQRS:MySQL 写 + Outbox + CDC 投影器 + Eager Invalidation + 乐观锁版本号 | | ai | 有状态工作流 | 无状态推理(AI 耗时工作流用 Temporal,状态在 Redis/Temporal DB) | | 配置 | env + DB 散落 | config-service 统一 | | 服务间通信 | gRPC + GraphQL 混用 | 外部 GraphQL(经 Router)+ 内部 gRPC(禁止子图互调) | | 服务数(业务+网关+前端) | 15 个(1 gw + 1 push + 3 BFF + 6 services + 4 portal) | 12 个(1 gw + 1 realtime + 1 router + 1 temporal + 7 services + 1 shell) | | 容器数(含数据层 6 个 + Temporal 1 个) | 21 | 18 |