- 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 列表传递
74 KiB
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 架构影响地图 v2.0(当前态)
- Portal Shell 插件化仪表盘设计 v2.1
- 0010 架构蓝图
- 项目规则
1. 背景与目标
1.1 v1 架构痛点诊断
当前 Edu 微服务架构(v1/v2.0 of 004)存在以下核心痛点:
- BFF 样板代码灾难:3 个手写 BFF(teacher-bff / student-bff / parent-bff)重复实现 GraphQL Resolver + gRPC 调用,每个字段改动需改 3 处(proto → schema → resolver),×3 灾难
- DataScope 跨库 JOIN 悖论:core-edu 要过滤"管理员可见班级",但用户-学校关系在 iam 库,微服务禁止跨库 JOIN,运行时传 ID 列表性能崩溃
- CDC 与 Outbox 文档矛盾:§1.1a 画 CDC(MySQL→Debezium→Kafka),§7.1 画 Outbox(业务表→Relay Worker→Kafka),看起来像两条路径做同一件事
- WebSocket 性能开销:push-gateway 同时支持 WS/SSE,但 K12 场景以单向通知为主,WS 心跳/握手开销过大
- Temporal 滥用边界:v1 规划引入 Temporal 但未明确边界,短事务若走 Temporal 会被入库开销拖垮吞吐;v2 严格立规矩——仅限 AI 耗时工作流 + Saga 分布式事务补偿
- iam 职责过载:iam 承载认证 + RBAC + 审计 + JWKS + 插件配置(6 张表)+ DataScope,成为"上帝服务"
- content 多存储耦合:content 同时写 MySQL + Neo4j + ES,事务一致性、同步延迟、故障域都是问题
- ai 有状态化:ai 服务内嵌 WorkflowStateStore,边界不清,难以水平扩展
- 配置散落:配置散落在 env / DB / code,无统一配置中心
- 前端 4 portal 重复:4 个独立 Next.js portal 重复实现 Dockerfile / i18n / auth / layout(已在 portal-shell v2.1 设计稿中解决)
1.2 v2 设计目标
- BFF 联邦化:采用 Apollo Federation,消灭手写 BFF 样板代码
- DataScope 运行时解析:通过 @requires 指令在子图间传递可见范围,解决跨库悖论
- CDC/Outbox 职责分工:Outbox 负责领域事件,CDC 负责读模型投影
- SSE 优先:推送默认走 SSE,WS 仅限强双向场景
- Temporal 严格边界:仅限 AI 耗时工作流(生成大纲→知识点→题目→组装试卷)+ 跨服务 Saga 分布式事务补偿(如购买插件:扣减积分服务 + 授权服务);CRUD 短事务绝对禁止
- 服务边界清晰:iam 拆分、content CQRS、ai 无状态化
- 配置中心化:config-service 统一管理动态配置
- 单 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 核心数据流
graph TB
User[用户] --> PortalShell[portal-shell :4010<br/>RSC 预取]
PortalShell --> ApiGateway[api-gateway :8080<br/>JWT + 限流]
ApiGateway --> ApolloRouter[apollo-router :4011<br/>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 扩展:
// 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 子图):
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):
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 处):
- 改
core_edu.proto,加int32 duration = 5;(手动) - 重新生成 proto 代码(自动)
- 改 teacher-bff 的
schema.graphql+ resolver(手动) - 改 student-bff 的
schema.graphql+ resolver(手动) - 改 parent-bff 的
schema.graphql+ resolver(手动) - 改 portal 的 GraphQL 查询(手动)
v2 流程(手动改动点 2 处):
- 改
core_edu.proto,加int32 duration = 5;(手动) - 运行
buf generate,自动生成 GraphQL schema(自动) - core-edu 服务自动暴露
duration字段(自动) - Apollo Router 自动拉取新 schema(自动)
- 改 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 解析器):
import DataLoader from "dataloader";
// 批量查询函数:接收 examId 数组,返回 Exam 数组
async function batchLoadExams(
examIds: readonly string[],
examRepo: ExamRepository,
): Promise<Exam[]> {
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<string, Exam>;
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<Exam> {
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 的查询计划与权限聚合,直接调子图。
实现:
# 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"
// 子图侧 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 子图设计:
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.recordedtopic) - 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:
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)
@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<string> {
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)
@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<Grade[]> {
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.<domain>.<aggregate>.<action> 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):
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):
{
"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 默认按 <topic.prefix>.<table> 输出,用 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 ← 兜底清理(防止网络抖动丢失)
代码示例:
@Injectable()
export class QuestionService {
constructor(
private questionRepo: QuestionRepository,
@Inject("REDIS") private redis: Redis,
private outbox: OutboxService,
) {}
async updateQuestion(id: string, dto: UpdateQuestionDto): Promise<Question> {
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 扩展:
type Question {
id: ID!
content: String
version: Int!
updatedAt: DateTime!
}
input QuestionQueryInput {
id: ID!
"""
客户端持有的最新版本号,用于读后一致性校验
"""
expectedVersion: Int
}
type Query {
question(input: QuestionQueryInput!): Question
}
Resolver 实现:
@Query(() => Question)
async question(
@Args("input") input: QuestionQueryInput,
): Promise<Question> {
// 先读 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.<domain>.<aggregate>.<action>
| 域 | 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 ChannelGET /ws— WebSocket 升级(监考专用)GET /online/:userId— 查询在线状态(检查 Redis SETonline: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 官方文档
- Apollo Router 配置
- GraphQL @requires 指令
- CQRS 模式
- Outbox 模式
- SSE vs WebSocket
- Portal Shell v2.1 设计稿
- 004 架构影响地图
- 项目规则
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 |