docs(docs): v2.1 架构 spec 8 项调整
- 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 列表传递
This commit is contained in:
@@ -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,21 +147,22 @@ 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)]
|
||||
```
|
||||
@@ -168,14 +170,14 @@ graph TB
|
||||
### 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 | 短事务禁用,长工作流轻量替代 |
|
||||
| Temporal | 规划引入但无边界 | 严格边界引入(AI 工作流 + Saga),CRUD 短事务禁用 | 长工作流自动重试/休眠/状态持久化,短事务走同步调用 |
|
||||
| iam 职责 | 认证+RBAC+插件配置 | 拆分 iam + config-service | 避免上帝服务 |
|
||||
| content 存储 | MySQL+Neo4j+ES 混写 | MySQL 写 + Outbox 投影器读模型 | CQRS 分离 |
|
||||
| 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<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 的查询计划与权限聚合,直接调子图。
|
||||
|
||||
**实现**:
|
||||
|
||||
```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<string[]> {
|
||||
if (ctx.dataScope === "ALL") return [];
|
||||
|
||||
const cacheKey = `iam:datascope:class:${ctx.userId}`;
|
||||
const cached = await this.redis.get(cacheKey);
|
||||
if (cached) return JSON.parse(cached);
|
||||
): 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);
|
||||
await this.redis.setex(cacheKey, 300, JSON.stringify(classIds));
|
||||
return classIds;
|
||||
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<Grade[]> {
|
||||
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.<domain>.<aggregate>.<action> 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 默认按 `<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/`)
|
||||
@@ -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<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 扩展**:
|
||||
|
||||
```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<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 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 命名规范
|
||||
|
||||
@@ -706,9 +1071,9 @@ content/src/
|
||||
### 6.4 幂等性
|
||||
|
||||
| 组件 | 幂等机制 |
|
||||
| ---------------- | ------------------------------------------------- |
|
||||
| ---------------- | ----------------------------------------------------------------------------------------------------- |
|
||||
| Kafka Producer | `idempotent=true` + `transactionalId` |
|
||||
| Outbox Publisher | 基于 `event_id` 去重(Redis SETNX) |
|
||||
| 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 事件 | 备课工作流、报告生成 |
|
||||
| 短事务 | 毫秒~秒 | 同步调用 + 分布式锁 | 作业提交、成绩录入 |
|
||||
| ------------------- | --------- | -------------------------------------------- | ------------------------------------------------------ |
|
||||
| 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 保留)
|
||||
|
||||
### 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 指令运行时解析 |
|
||||
| **Outbox 领域事件** | 领域事件必须走 Outbox,禁止直接调 Kafka producer |
|
||||
| **CDC 仅备选投影** | CDC/Debezium 默认不启用,仅当 Outbox 事件粒度不够时作为投影器备选;禁止用于领域事件发布 |
|
||||
| **DataScope @requires** | 禁止跨库 JOIN,DataScope 通过 @requires 指令运行时解析(v2.1 用 ScopeToken 优化) |
|
||||
| **Outbox 领域事件** | 领域事件必须走 Outbox 表,禁止直接调 Kafka producer |
|
||||
| **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 |
|
||||
| **ai 无状态** | ai 不存业务数据,会话/工作流状态存 Redis(AI 耗时工作流用 Temporal) |
|
||||
| **单 portal-shell** | 前端单容器部署,禁止新增独立 portal |
|
||||
|
||||
**保留约束(v1)**:
|
||||
### 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 |
|
||||
| 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 实现 | 教师查询可见班级成绩正确 | M2 |
|
||||
| M5 | content CQRS 改造(投影器) | Neo4j/ES 通过投影器同步 | M1 |
|
||||
| M4 | DataScope @requires + ScopeToken 实现 | 教师查询可见班级成绩正确 | M2 |
|
||||
| M5 | content CQRS 改造(投影器 + Eager Invalidation + 乐观锁) | Neo4j/ES 通过投影器同步,写后读一致 | M0.5/M1 |
|
||||
| M6 | ai 无状态化 | ai 不存业务数据,状态在 Redis | M1 |
|
||||
| M7 | realtime-gateway SSE 优先 | SSE 推送可用 | 无 |
|
||||
| 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 | arch.db + 004 文档同步 | arch:scan 通过 | 全部 |
|
||||
| M11 | 废弃各服务 OutboxPublisher 轮询线程 | 无服务使用轮询投递 | M0.5 |
|
||||
| M12 | arch.db + 004 文档同步 | arch:scan 通过 | 全部 |
|
||||
|
||||
### 10.2 并行运行策略
|
||||
|
||||
@@ -896,7 +1314,7 @@ 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 |
|
||||
@@ -904,11 +1322,18 @@ realtime-gateway 消费 Kafka(edu.notify.notification.*)
|
||||
| ADR-027 | content CQRS 改造 | 写读分离,故障隔离 |
|
||||
| ADR-028 | ai 无状态化 | 明确边界,可扩展 |
|
||||
| ADR-029 | SSE 优先 + WS 仅监考 | K12 场景单向为主 |
|
||||
| ADR-030 | 不引入 Temporal | 短事务禁用,长工作流用轻量状态机 |
|
||||
| ADR-030 | Temporal 严格边界引入(v2.1 修订) | AI 工作流 + Saga 用 Temporal;CRUD 短事务绝对禁止,走同步调用 + 分布式锁 |
|
||||
| ADR-031 | config-service + Redis 轻量配置中心 | K12 场景,不引入 etcd |
|
||||
| ADR-032 | CDC 默认禁用,投影器首选 Outbox 事件 | 与 Outbox 职责分工,简化依赖 |
|
||||
| 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,19 +1345,20 @@ 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/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 子图 | 改造 |
|
||||
| `services/push-gateway/` | 改名 realtime-gateway;SSE 优先 | 改造 |
|
||||
| `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 下线服务
|
||||
@@ -950,17 +1376,21 @@ 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-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(本设计) |
|
||||
| -------------------------------------------------------------- | ------------------------------------------------------ | -------------------------------------------------------------------------------------------------------------------- |
|
||||
| 维度 | v1(当前) | v2.1(本设计) |
|
||||
| --------------------------------------- | ------------------------------------------------------ | ------------------------------------------------------------------------- |
|
||||
| 前端 | 4 portal + MF | 单 portal-shell |
|
||||
| BFF | 3 个手写 BFF | Apollo Federation |
|
||||
| DataScope | 网关透传 + 服务注入 | @requires 运行时解析 |
|
||||
| Push | WS + SSE | SSE 优先 + WS 仅监考 |
|
||||
| Temporal | 规划引入 | 不引入 |
|
||||
| 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 写 + 投影器 |
|
||||
| ai | 有状态工作流 | 无状态推理 |
|
||||
| content | MySQL+Neo4j+ES 混写 | CQRS:MySQL 写 + Outbox + CDC 投影器 + Eager Invalidation + 乐观锁版本号 |
|
||||
| ai | 有状态工作流 | 无状态推理(AI 耗时工作流用 Temporal,状态在 Redis/Temporal DB) |
|
||||
| 配置 | 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 |
|
||||
| 服务间通信 | 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 |
|
||||
|
||||
Reference in New Issue
Block a user