From 754595620234e5440b64d373d43e4580a407536d Mon Sep 17 00:00:00 2001 From: SpecialX <47072643+wangxiner55@users.noreply.github.com> Date: Tue, 14 Jul 2026 22:55:56 +0800 Subject: [PATCH] =?UTF-8?q?docs(docs):=20v2.1=20=E6=9E=B6=E6=9E=84=20spec?= =?UTF-8?q?=208=20=E9=A1=B9=E8=B0=83=E6=95=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Temporal 严格边界引入(AI 工作流 + Saga,CRUD 短事务禁用) - CDC+Outbox 结合(Debezium Transaction Log Tailing,废弃轮询) - DataLoader 强制(@key 解析器请求合并,消除 N+1) - Router-Authorization 信任凭证(拒绝非 Router 的 GraphQL 请求) - 外部 GraphQL + 内部 gRPC 边界明确 - Eager Invalidation + 乐观锁版本号回传 - Redis Pub/Sub 推送背板(边缘网关不挂 Kafka) - ScopeToken 优化大规模 ID 列表传递 --- ...6-07-14-architecture-v2-redesign-design.md | 806 ++++++++++++++---- 1 file changed, 625 insertions(+), 181 deletions(-) diff --git a/docs/superpowers/specs/2026-07-14-architecture-v2-redesign-design.md b/docs/superpowers/specs/2026-07-14-architecture-v2-redesign-design.md index 8199f25..547843b 100644 --- a/docs/superpowers/specs/2026-07-14-architecture-v2-redesign-design.md +++ b/docs/superpowers/specs/2026-07-14-architecture-v2-redesign-design.md @@ -1,8 +1,9 @@ # Edu 平台架构 v2 重设计 -> 版本:v2.0 +> 版本: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(当前态) @@ -22,7 +23,7 @@ 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 滥用边界**:规划引入 Temporal 但未明确边界,短事务若走 Temporal 会被入库开销拖垮吞吐 +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,边界不清,难以水平扩展 @@ -35,7 +36,7 @@ 2. **DataScope 运行时解析**:通过 @requires 指令在子图间传递可见范围,解决跨库悖论 3. **CDC/Outbox 职责分工**:Outbox 负责领域事件,CDC 负责读模型投影 4. **SSE 优先**:推送默认走 SSE,WS 仅限强双向场景 -5. **Temporal 不引入**:长工作流用 Redis 状态机 + Outbox 事件 +5. **Temporal 严格边界**:仅限 AI 耗时工作流(生成大纲→知识点→题目→组装试卷)+ 跨服务 Saga 分布式事务补偿(如购买插件:扣减积分服务 + 授权服务);CRUD 短事务绝对禁止 6. **服务边界清晰**:iam 拆分、content CQRS、ai 无状态化 7. **配置中心化**:config-service 统一管理动态配置 8. **单 portal-shell**:前端单容器部署(引用 portal-shell v2.1 设计稿) @@ -46,14 +47,14 @@ - **技术栈全开放**:可换 DB / MQ / 框架 / 语言 - **数据可清零**:无迁移包袱,允许重建测试数据 - **微服务架构保留**:不为当前规模妥协,为未来扩展性预留 -- **AI 禁止切换分支**:分支创建/切换由人类决策者负责 +- **AI 分支权限**:v2.1 实施期间,人类决策者已授予 AI 提交其他模块代码与切换分支的权限(仅限 v2 实施期,结束后回收) ### 1.4 非目标 - 不重写已有业务逻辑(仅调整架构边界) - 不替换 MySQL / Redis / Kafka / ClickHouse / Neo4j / Elasticsearch - 不引入 etcd / Consul(K12 场景过重) -- 不引入 Temporal / Cadence(短事务禁用,长工作流用轻量状态机) +- 不将 Temporal 用于 CRUD 短事务(作业提交、成绩录入等毫秒~秒级事务走同步调用 + 分布式锁) --- @@ -123,8 +124,8 @@ L6 数据层 **移除的组件**: -- Debezium Connect → v2 默认不启用,投影器走 Outbox 事件(见 §5.4 / §6.1);仅当 Outbox 事件无法覆盖投影需求时,作为备选方案启用 -- Temporal → 不引入(ai 保留 WorkflowStateStore 作为轻量替代) +- Debezium Connect → 保留并启用,与 Outbox 结合使用(见 §5.4 / §6.1):业务代码只写业务表 + Outbox 表,去掉 OutboxPublisher 轮询线程,改由 Debezium 监听 binlog 自动将 Outbox 表变更推送到 Kafka(Transaction Log Tailing) +- Temporal → 引入但严格立规矩(见 §8.2):仅限 AI 耗时工作流 + Saga 分布式事务补偿,CRUD 短事务绝对禁止 ### 2.3 核心数据流 @@ -146,38 +147,39 @@ graph TB CoreEdu -.接收 visibleClassIds.-> CoreEdu PortalShell -.SSE 推送.-> RealtimeGw[realtime-gateway :8081] - RealtimeGw -.Kafka 消费.-> Msg + RealtimeGw -.Redis Pub/Sub 订阅.-> Msg + Msg -.Kafka 消费.-> Kafka - Iam <-->|Outbox| Kafka[(Kafka)] - CoreEdu <-->|Outbox| Kafka - Content <-->|Outbox + 投影器消费| Kafka - Msg <-->|Outbox| Kafka - DataAna <-->|Outbox 消费| Kafka - Ai <-->|Outbox| 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 -.Outbox 投影器.-> Neo4j[(Neo4j)] - Content -.Outbox 投影器.-> ES[(Elasticsearch)] + 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 | 规划引入 | 不引入,保留 WorkflowStateStore | 短事务禁用,长工作流轻量替代 | -| iam 职责 | 认证+RBAC+插件配置 | 拆分 iam + config-service | 避免上帝服务 | -| content 存储 | MySQL+Neo4j+ES 混写 | MySQL 写 + Outbox 投影器读模型 | CQRS 分离 | -| ai 状态 | 有状态工作流 | 无状态推理 + WorkflowStateStore | 明确边界 | -| 配置 | env + DB 散落 | config-service 统一 | 配置中心化 | +| 决策点 | 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 统一 | 配置中心化 | --- @@ -331,6 +333,140 @@ headers: **手动改动点从 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. 业务服务层调整 @@ -496,28 +632,37 @@ query GetGradesForMyClasses($termId: ID!) { } ``` -**Apollo Router 查询计划**(自动生成): +**Apollo Router 查询计划**(自动生成,使用 ScopeToken 优化): ``` 步骤 1: 调 iam 子图 - visibleClassIds(userId=x-user-id) + scopeToken(userId=x-user-id, scopeType="class") → iam 查 Redis 缓存(key: iam:datascope:class:{userId},TTL 5min) → 未命中则查 user_school_role 表 - → 返回 ["class-001", ..., "class-050"] + → 将 ID 列表存入 Redis Set(key: scope:usr:{userId}:cls_scope,TTL 5min) + → 返回极短的 scopeToken: "usr:123:cls_scope"(不传全量 ID 数组) 步骤 2: 调 core-edu 子图 gradesForCurrentUser( - visibleClassIds: ["class-001", ..., "class-050"], + scopeToken: "usr:123:cls_scope", termId: "2024-spring" ) - @requires(fields: "visibleClassIds") - → core-edu 查 MySQL: WHERE class_id IN (...) AND term_id = ? + @requires(fields: "scopeToken") + → core-edu 从同机房 Redis SMEMBERS scope:usr:123:cls_scope 拿到 ID 数组 + → 查 MySQL: WHERE class_id IN (...) AND term_id = ? → 返回 [Grade, Grade, ...] 步骤 3: 聚合返回客户端 ``` -#### iam 子图 DataScope Resolver +**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) @@ -527,30 +672,39 @@ export class DataScopeResolver { @Inject("REDIS") private redis: Redis, ) {} - @ResolveField(() => [String]) - async visibleClassIds( + @ResolveField(() => String) + async classScopeToken( @Parent() user: User, @Context() ctx: { userId: string; dataScope: string }, - ): Promise { - if (ctx.dataScope === "ALL") return []; + ): Promise { + if (ctx.dataScope === "ALL") return "ALL"; - const cacheKey = `iam:datascope:class:${ctx.userId}`; - const cached = await this.redis.get(cacheKey); - if (cached) return JSON.parse(cached); - - const classIds = await this.iamService.getVisibleClassIds(ctx.userId); - await this.redis.setex(cacheKey, 300, JSON.stringify(classIds)); - return classIds; + 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 +#### core-edu 子图 @requires Resolver(从 ScopeToken 解析 ID) ```typescript @Resolver(() => Grade) export class GradeResolver { - constructor(private gradeService: GradeService) {} + constructor( + private gradeService: GradeService, + @Inject("REDIS") private redis: Redis, + ) {} @Query(() => [Grade]) @RequirePermission("GRADES_READ") @@ -559,14 +713,17 @@ export class GradeResolver { ctx: { userId: string; dataScope: string; - visibleClassIds?: string[]; + classScopeToken?: string; }, @Args("termId") termId: string, ): Promise { - if (ctx.dataScope === "ALL") { + if (ctx.dataScope === "ALL" || ctx.classScopeToken === "ALL") { return this.gradeService.findByTerm(termId); } - return this.gradeService.findByClassIds(ctx.visibleClassIds ?? [], termId); + // 从 Redis Set 批量获取实际 ID + const redisKey = `scope:${ctx.classScopeToken}`; + const classIds = await this.redis.smembers(redisKey); + return this.gradeService.findByClassIds(classIds, termId); } } ``` @@ -600,27 +757,80 @@ config-service 发 Kafka 事件(edu.config.entry.changed) - 特性开关(feature_flags,如 `enable_ai_tutor`) - 系统配置(system_config,如 `max_exam_duration`) -### 5.4 Outbox 投影器设计(content CQRS) +### 5.4 Outbox + CDC 结合(Debezium Transaction Log Tailing) -**投影器架构**(默认走 Outbox 事件;CDC/Debezium 仅作备选,当 Outbox 事件粒度不够时启用): +**核心原则**:业务代码依然只写业务表 + Outbox 表(保证业务语义),但去掉 OutboxPublisher 轮询线程,改由 Debezium 监听 binlog 自动将 Outbox 表变更推送到 Kafka(Transaction Log Tailing 模式)。 + +**架构**: ``` -content 服务写 MySQL - ↓ Outbox 事件 +写侧(业务服务): + 业务事务内写业务表 + outbox 表(原子) + ↓ 不再有轮询线程 ↓ -Kafka(edu.content.* topic) +MySQL binlog(row mode) + ↓ Debezium Connect 监听(伪装成 MySQL Slave) ↓ - ├── neo4j-projector(消费 → 写 Neo4j 知识图谱) - │ · 监听 KnowledgePointCreated / KnowledgePointLinked - │ · 写 Neo4j 节点 + 关系 +Kafka(edu... topic) + ↓ + ├── 投影器(CQRS 读模型同步) + │ ├── neo4j-projector(content → Neo4j 知识图谱) + │ ├── es-projector(content → ES 题库索引) + │ ├── content-cache-projector(content → Redis 缓存失效,兜底) + │ └── msg-cache-projector(msg → Redis 缓存失效,兜底) │ - ├── es-projector(消费 → 写 ES 题库索引) - │ · 监听 QuestionCreated / QuestionUpdated / QuestionDeleted - │ · 写 ES 索引(题干、选项、知识点标签) - │ - └── content-cache-projector(消费 → 失效 Redis 缓存) - · 监听所有 content 事件 - · 失效 content 查询缓存 + └── 业务消费者(领域事件处理) + ├── 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/`) @@ -637,32 +847,187 @@ content/src/ └── 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 vs CDC 职责分工 +### 6.1 Outbox + CDC 结合模式(Transaction Log Tailing) + +**核心决策**:v2.1 不再在 Outbox 和 CDC 之间二选一,而是**结合使用**——业务代码只写 Outbox 表(业务语义),由 Debezium 监听 binlog 自动投递到 Kafka(技术传输)。 ``` -写侧:Outbox 模式(领域事件发布) +写侧:业务表 + Outbox 表(同事务原子) 业务服务(iam/core-edu/content/msg/ai) ① 业务事务内写业务表 + outbox 表(原子) - ② OutboxPublisher 轮询 outbox 表 → 投递 Kafka - ③ 保证 at-least-once 语义 - 用途:发布业务语义事件(ExamCreated / HomeworkSubmitted) + ② 不再有 OutboxPublisher 轮询线程(v1 的轮询模式被废弃) + ③ Debezium Connect 监听 MySQL binlog(伪装成 Slave) + ④ Debezium 监听到 outbox 表有新增 → 自动推送到 Kafka + ⑤ 保证 at-least-once 语义(Debezium offset 管理) + 优势:业务代码不感知 MQ;无损监听 binlog 不影响业务性能;消除轮询线程的延迟和 DB 压力 -读侧:投影器(CQRS 读模型 / 跨服务物化视图) - 方式 A:Debezium → Kafka(用于跨服务物化视图同步,备选) - 方式 B:Outbox 事件 → 投影器(用于读模型投影,推荐) +读侧:投影器消费 Kafka(CQRS 读模型同步) + ├── neo4j-projector(content → Neo4j 知识图谱) + ├── es-projector(content → ES 题库索引) + ├── content-cache-projector(content → Redis 缓存失效,兜底) + └── msg-cache-projector(msg → Redis 缓存失效,兜底) 用途:CQRS 读模型投影,非业务事件发布 ``` **关键原则**: -- 领域事件必须走 Outbox(业务语义,由业务代码显式发布) -- 数据投影走 CDC 或 Outbox 消费(技术同步,无需业务语义) -- CDC 不替代 Outbox(CDC 是 binlog 级,Outbox 是业务语义级) -- v2 简化:优先用 Outbox 事件做投影器,减少 Debezium 依赖 +- 业务代码只写业务表 + Outbox 表(业务语义由业务代码控制) +- Outbox 表的投递由 Debezium 自动完成(技术传输由基础设施完成) +- 投影器消费 Kafka 事件做 CQRS 读模型同步(与业务事件发布解耦) +- Outbox + CDC 不矛盾,而是分工:Outbox 定义"发什么",CDC 解决"怎么发" +- v1 的 OutboxPublisher 轮询线程在 v2.1 完全废弃(被 Debezium 替代) ### 6.2 Kafka Topic 命名规范 @@ -705,12 +1070,12 @@ content/src/ ### 6.4 幂等性 -| 组件 | 幂等机制 | -| ---------------- | ------------------------------------------------- | -| Kafka Producer | `idempotent=true` + `transactionalId` | -| Outbox Publisher | 基于 `event_id` 去重(Redis SETNX) | -| Consumer | 基于 `event_id` 去重(DB 唯一索引或 Redis SETNX) | -| 投影器 | 基于 `event_id` + `aggregate_id` 去重 | +| 组件 | 幂等机制 | +| ---------------- | ----------------------------------------------------------------------------------------------------- | +| 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` 去重 | --- @@ -774,39 +1139,61 @@ core-edu 子图 @requires 接收 visibleClassIds,注入 WHERE - 内置 OpenTelemetry trace - 内置持续查询(Persisted Queries)统计 -### 8.2 工作流边界(Temporal 约束) +### 8.2 工作流边界(Temporal 严格约束) -**v2 决策**:不引入 Temporal,保留 ai 服务的 WorkflowStateStore 作为轻量替代。 +**v2.1 决策**:引入 Temporal,但严格立规矩——仅限 AI 耗时工作流 + Saga 分布式事务补偿,CRUD 短事务绝对禁止。 **架构约束**: -| 工作流类型 | 时长 | 推荐方案 | 示例 | -| ---------- | ------- | -------------------------- | -------------------- | -| 长工作流 | 分钟~天 | Redis 状态机 + Outbox 事件 | 备课工作流、报告生成 | -| 短事务 | 毫秒~秒 | 同步调用 + 分布式锁 | 作业提交、成绩录入 | -| 批处理 | 小时 | Cron + Batch Job | 成绩统计、掌握度计算 | +| 工作流类型 | 时长 | 推荐方案 | 示例 | +| ------------------- | --------- | -------------------------------------------- | ------------------------------------------------------ | +| AI 耗时工作流 | 分钟~小时 | **Temporal Workflow** | 生成大纲 → 生成知识点 → 生成题目 → 组装试卷 | +| Saga 分布式事务补偿 | 秒~分钟 | **Temporal Saga** | 购买高级插件:扣减积分服务(+补偿)+ 授权服务(+补偿) | +| 长工作流(非 AI) | 分钟~天 | Redis 状态机 + Outbox 事件 | 备课工作流(非 AI 部分)、报告生成 | +| 短事务 | 毫秒~秒 | **同步调用 + 分布式锁**(绝对禁止 Temporal) | 作业提交、成绩录入、考勤记录 | +| 批处理 | 小时 | Cron + Batch Job | 成绩统计、掌握度计算 | **禁止**: -- 短事务用工作流引擎(入库开销拖垮吞吐) +- **CRUD 短事务绝对禁止用 Temporal**:作业提交、成绩录入、考勤记录等毫秒~秒级事务走同步调用 + 分布式锁,Temporal 的入库开销会拖垮吞吐量 - 纯读操作走工作流(直接走缓存) **允许**: -- 跨天长流程用状态机 + 事件驱动 -- ai 备课/报告生成用 Redis 状态机(WorkflowStateStore) +- AI 耗时工作流用 Temporal(自动重试 / 休眠 / 状态持久化) +- 跨服务 Saga 分布式事务补偿用 Temporal(自动补偿 / 状态可查询) +- 跨天长流程用 Redis 状态机 + Outbox 事件(非 AI 场景) -### 8.3 推送架构(SSE 优先) +**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 → Kafka + ↓ Outbox 表 + Debezium → Kafka(edu.notify.notification.*) ↓ -realtime-gateway 消费 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,保持长连接 - │ · 服务端 push 事件流 + │ · realtime-gateway 收到 Redis Pub/Sub 消息 → 推送到 SSE 连接 │ · 适合 99% 场景(考试发布/成绩推送/通知) │ └── WebSocket 推送(仅监考场景) @@ -815,43 +1202,71 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*) · 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 推送(默认) +- `GET /sse` — SSE 推送(默认),连接时 SUBSCRIBE Redis Channel - `GET /ws` — WebSocket 升级(监考专用) -- `POST /internal/push` — msg 服务 HTTP 调用推送 -- `GET /online/:userId` — 查询在线状态 +- `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 新增) +## 9. 架构约束(v2.1 新增 + v2 保留) -| 约束 | 说明 | -| ----------------------- | --------------------------------------------------------------------------------------- | -| **BFF 联邦化** | 禁止手写 BFF 聚合层,所有 GraphQL 查询走 Apollo Router | -| **子图自主** | 每个业务服务必须暴露 GraphQL 子图,proto → GraphQL 自动生成 | -| **DataScope @requires** | 禁止跨库 JOIN,DataScope 通过 @requires 指令运行时解析 | -| **Outbox 领域事件** | 领域事件必须走 Outbox,禁止直接调 Kafka producer | -| **CDC 仅备选投影** | CDC/Debezium 默认不启用,仅当 Outbox 事件粒度不够时作为投影器备选;禁止用于领域事件发布 | -| **SSE 优先** | 推送默认走 SSE,WS 仅限强双向场景(监考) | -| **Temporal 禁用** | 不引入 Temporal,长工作流用 Redis 状态机 + Outbox | -| **iam 职责限定** | iam 仅负责认证/RBAC/JWKS/审计/DataScope,禁止承载配置 | -| **config 统一** | 动态配置必须走 config-service,禁止散落 env/code | -| **content CQRS** | content 仅写 MySQL,Neo4j/ES 通过投影器同步 | -| **ai 无状态** | ai 不存业务数据,会话/工作流状态存 Redis | -| **单 portal-shell** | 前端单容器部署,禁止新增独立 portal | +### 9.1 v2.1 新增约束 -**保留约束(v1)**: +| 约束 | 说明 | +| ---------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------- | +| **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(事件) +- 服务间通信:GraphQL(@requires,仅前端)+ gRPC(服务间)+ Kafka(事件) - 每服务独占 DB schema - JWT RS256 + JWKS - Cookie: httpOnly + Secure + SameSite=Strict @@ -861,22 +1276,25 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*) ## 10. 迁移路径 -### 10.1 阶段划分(从 v1 到 v2) +### 10.1 阶段划分(从 v1 到 v2.1) -| 阶段 | 内容 | 验收标准 | 依赖 | -| ---- | ----------------------------------------- | ---------------------------------- | ----- | -| M0 | proto → GraphQL 代码生成工具链 | `buf generate` 输出 GraphQL schema | 无 | -| M1 | 各服务暴露 GraphQL 子图 | `/graphql` 端点可查询 | M0 | -| M2 | apollo-router 部署 + Supergraph 组装 | Router 可聚合查询 | M1 | -| M3 | iam 拆分 config-service | config-service 独立运行 | M1 | -| M4 | DataScope @requires 实现 | 教师查询可见班级成绩正确 | M2 | -| M5 | content CQRS 改造(投影器) | Neo4j/ES 通过投影器同步 | M1 | -| M6 | ai 无状态化 | ai 不存业务数据,状态在 Redis | M1 | -| M7 | realtime-gateway SSE 优先 | SSE 推送可用 | 无 | -| M8 | portal-shell 接入 apollo-router | portal-shell 查询走 Router | M2/M4 | -| M9 | 旧 BFF 下线(teacher/student/parent-bff) | 流量为 0 | M8 | -| M10 | 旧 portal 下线 | 流量为 0 | M8 | -| M11 | arch.db + 004 文档同步 | arch:scan 通过 | 全部 | +| 阶段 | 内容 | 验收标准 | 依赖 | +| ---- | --------------------------------------------------------- | ----------------------------------------- | ------- | +| 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 并行运行策略 @@ -895,20 +1313,27 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*) ## 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 | 短事务禁用,长工作流用轻量状态机 | -| ADR-031 | config-service + Redis 轻量配置中心 | K12 场景,不引入 etcd | -| ADR-032 | CDC 默认禁用,投影器首选 Outbox 事件 | 与 Outbox 职责分工,简化依赖 | -| ADR-033 | portal-shell 单容器替代 4 portal | Modular Monolith | -| ADR-034 | realtime-gateway 改名 | SSE 优先定位 | +| 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 开销 | --- @@ -920,20 +1345,21 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*) | -------------------------- | -------------------------------------------------- | ---- | | `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 子图 | 改造 | -| `services/core-edu/` | 新增 GraphQL 子图;实现 @requires DataScope | 改造 | -| `services/content/` | CQRS 改造:移除直接写 Neo4j/ES,新增投影器 worker | 改造 | -| `services/msg/` | 新增 GraphQL 子图 | 改造 | -| `services/data-ana/` | 新增 GraphQL 子图 | 改造 | -| `services/ai/` | 无状态化:WorkflowStateStore 改为 Redis;新增 GraphQL 子图 | 改造 | -| `services/push-gateway/` | 改名 realtime-gateway;SSE 优先 | 改造 | -| `services/api-gateway/` | 移除 BFF 代理路由;新增 apollo-router 路由 | 改造 | +| 路径 | 变更 | 动作 | +| ------------------------ | ---------------------------------------------------------------------------------------------- | ---- | +| `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 下线服务 @@ -949,18 +1375,22 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*) ### 12.4 新增工具链 -| 路径 | 职责 | 动作 | -| -------------------------------------- | ----------------------------- | ---- | -| `packages/shared-proto/gen-graphql.ts` | proto → GraphQL schema 生成器 | 新建 | -| `packages/shared-proto/buf.gen.yaml` | 新增 GraphQL 输出配置 | 修改 | +| 路径 | 职责 | 动作 | +| -------------------------------------- | ------------------------------------------- | ---- | +| `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 / portal-shell;移除 3 BFF + 4 portal | 修改 | -| `infra/port-allocation.md` | 新增 3011 / 4000 / 4010 端口 | 修改 | -| `infra/init-sql/` | 新增 config-service schema DDL | 新增 | +| 路径 | 变更 | 动作 | +| --------------------------------- | ------------------------------------------------------------------------------------ | ---- | +| `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 镜像预拉 | 修改 | --- @@ -969,11 +1399,20 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*) ### 13.1 功能验收 - [ ] Apollo Router 可聚合 7 个子图查询 -- [ ] DataScope @requires 正确传递 visibleClassIds +- [ ] DataScope @requires + ScopeToken 正确传递可见范围 - [ ] iam 拆分后 config-service 独立运行 - [ ] content CQRS 投影器正确同步 Neo4j/ES +- [ ] Eager Invalidation:写后同步 Redis DEL 生效 +- [ ] 乐观锁版本号:写后读一致(ES 落后则穿透读 MySQL) - [ ] ai 无状态化,工作流状态在 Redis -- [ ] realtime-gateway SSE 推送可用 +- [ ] 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 事件失效各服务缓存 @@ -981,8 +1420,11 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*) - [ ] `pnpm run lint` + `pnpm run typecheck` 零错误 - [ ] Apollo Router 查询延迟 P99 < 100ms(不含子图) -- [ ] DataScope @requires 查询延迟 P99 < 500ms(含子图) +- [ ] DataScope @requires + ScopeToken 查询延迟 P99 < 500ms(含子图) +- [ ] DataLoader 批量加载:并发 100 次 @key 解析,Repository 调用 ≤ 1 次 - [ ] SSE 单节点可扛 10 万连接 +- [ ] Redis Pub/Sub 推送延迟 < 1ms +- [ ] Debezium 投递延迟 < 1s(binlog → Kafka) - [ ] arch.db 更新,004 文档同步 - [ ] 所有 ADR 记录完整 @@ -1004,16 +1446,18 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*) ## 15. 版本演进对比 -| 维度 | v1(当前) | v2(本设计) | -| -------------------------------------------------------------- | ------------------------------------------------------ | -------------------------------------------------------------------------------------------------------------------- | -| 前端 | 4 portal + MF | 单 portal-shell | -| BFF | 3 个手写 BFF | Apollo Federation | -| DataScope | 网关透传 + 服务注入 | @requires 运行时解析 | -| Push | WS + SSE | SSE 优先 + WS 仅监考 | -| Temporal | 规划引入 | 不引入 | -| iam | 认证+RBAC+插件配置 | 拆分 iam + config-service | -| content | MySQL+Neo4j+ES 混写 | CQRS:MySQL 写 + 投影器 | -| ai | 有状态工作流 | 无状态推理 | -| 配置 | env + DB 散落 | config-service 统一 | -| 服务数(业务+网关+前端) | 15 个(1 gw + 1 push + 3 BFF + 6 services + 4 portal) | 11 个(1 gw + 1 realtime + 1 router + 7 services + 1 shell;替换 3 BFF + 4 portal 为 1 Router + 1 Shell + 1 config) | -| 容器数(含数据层 6 个:MySQL/Redis/Kafka/ClickHouse/Neo4j/ES) | 21 | 17 | +| 维度 | 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 |