# 模块架构设计文档 — data-ana
> AI 标识:ai11(v2 由 ai11 审核 ai06 v1 后修订)
> 负责模块:data-ana(P4)
> 阶段:架构设计外包 · 阶段 2(模块架构设计)v2 审核版
> 日期:2026-07-09(v1 by ai06)/ 2026-07-10(v2 审核修订 by ai11)
> 关联文档:[阶段 1 理解确认书](./01-understanding.md)、[ai-allocation.md](../../../docs/architecture/ai-allocation.md)、[004 架构影响地图](../../../docs/architecture/004_architecture_impact_map.md)、[pending-features.md](../../../docs/architecture/roadmap/pending-features.md)
> 审查状态:v1 已通过 coord 交叉审查(见 [coord-cross-review.md](../../../docs/architecture/coord-cross-review.md)),裁决结论已沉淀到 004 §15.3;v2 在 v1 基础上对齐已仲裁约束 + 补强长远性/全面性
> v2 修订要点:(1) 身份归属 ai06 → ai11;(2) 响应信封对齐 ActionState(004 §11.5 P0 整改);(3) 删除已仲裁的"未决设计决策";(4) 补 gRPC 端口 50055;(5) 补 5 张宽表(ai-allocation §5 要求的 attendance_logs);(6) 补 4 端 Dashboard + 预警 API;(7) 补掌握度计算算法具体公式;(8) 补缓存策略/CDC 水平扩展/容量规划/数据治理/测试策略/Dockerfile 多阶段/配置项清单/API 版本化/Projection 演进路径/未来阶段铺垫(§9-§19 新增章节)
---
## 设计原则与全局约束
本设计遵循以下强制约束(来自 project_rules.md + coding-standards.md + 004):
1. **契约先行**:proto 已定义(analytics.proto / events.proto / iam.proto),实现前不修改 proto,如需修改走 coord 流程(ai-allocation.md §9.4)
2. **CQRS 读写分离**:data-ana 是纯读模型服务(无 MySQL 写),ClickHouse 宽表由 CDC 投影构建
3. **事件驱动**:data-ana 消费 CDC(主通道)+ 领域事件(备通道,待 P4 后期评估);发布 `edu.insight.mastery.updated` 派生数据事件(豁免 Outbox,004 §12.2 + §15.3 #6 已仲裁)
4. **gRPC 优先**:004 §4.1 + §4.2 明确 BFF → 业务服务走 gRPC,**P4 启用 data-ana gRPC server 端口 50055**(HTTP 保留作 Gateway 直连降级)
5. **DataScope 过滤**:004 §5.3 DataScope 6 级在查询层注入 WHERE;iam `GetEffectiveDataScope` gRPC(004 §15.3 #5 已仲裁 P4 补全)
6. **三支柱可观测**:structlog + prometheus-client + OpenTelemetry(已具备,需补业务指标 + gRPC server interceptor)
7. **统一响应信封 ActionState**(004 §11.5 已仲裁 P0 整改):成功 `{success: true, data: T}` / 失败 `{success: false, error: {code, message, details?, traceId?}}` / 降级 `degraded` 作为 `details.degraded` 子字段
8. **降级模式**:外部依赖(ClickHouse / Kafka / iam gRPC / Redis)不可用时返回骨架数据 + `details.degraded: true`
9. **Python 规范**:pydantic-settings 配置 / Pydantic 模型校验 / async 优先 / 类型注解强制 / ruff 零错误
10. **长远架构演进**:为 P5(ai 用量消费 / gRPC stream)、P6(CDC 水平扩展 / 容量规划 / 数据治理 / Service Mesh)做好铺垫,见 §14 / §19
---
## 1. 模块内部分层图
```mermaid
flowchart TB
subgraph Entry["入口层"]
HTTP[FastAPI HTTP Router
:3006 /analytics/* + /healthz + /readyz]
GRPC[grpc.aio Server
:50055 AnalyticsService]
end
subgraph Middleware["中间件层"]
AUTH[AuthDepends
校验 x-user-id / x-user-roles]
SCOPE[DataScopeDepends
注入 data_scope 元数据]
TRACE[OTel FastAPIInstrumentor
+ grpc.aio server interceptor]
end
subgraph Service["应用服务层 Application Service"]
S1[AnalyticsService
班级/学生/趋势/Dashboard 查询编排]
S2[MasteryService
掌握度计算 + 事件发布]
S3[ErrorBookService
错题本查询]
S4[WarningService
预警阈值评估 + 触发]
end
subgraph Repo["数据访问层 Repository"]
R1[ClickHouseRepository
宽表查询 + DataScope WHERE 注入 + FINAL/argMax]
R2[KafkaProducer
mastery.updated 派生数据事件发布]
R3[IamClient
gRPC 调 iam.GetEffectiveDataScope]
R4[RedisClient
DataScope 缓存 + CDC 幂等 + 预警位图]
end
subgraph Consumer["CDC 消费者(后台任务)"]
C1[CdcConsumer
aiokafka AIOKafkaConsumer]
C2[ExamCache
exam_id→class_id 映射 P4 内存 / P6 Redis 化]
C3[EventHandler
grades/exams/homework/classes/attendance/kp 路由]
end
subgraph Storage["存储 / 总线"]
CH[(ClickHouse
edu_analytics 库 5 宽表)]
KAFKA[(Kafka
edu-cdc.* + edu.insight.mastery.updated + edu.insight.ai.usage)]
IAM[iam:3002 gRPC :50052]
REDIS[(Redis
data_ana:* 键前缀)]
end
HTTP --> AUTH --> SCOPE --> S1
HTTP --> S3
HTTP --> S4
GRPC --> S1
S1 --> R1
S3 --> R1
S4 --> R1
S4 --> R2
S2 --> R1
S2 --> R2
SCOPE --> R3
R3 --> R4
R1 --> CH
R2 --> KAFKA
R3 --> IAM
R4 --> REDIS
C1 --> C3
C3 --> C2
C3 --> R1
C3 --> R4
C1 --> KAFKA
```
**分层规则**:
- **入口层**:HTTP `:3006`(保留作 Gateway 直连降级)+ gRPC `:50055`(P4 启用主入口,BFF / ai 调用)。两入口共享同一 Application Service
- **中间件层**:FastAPI Depends 链(`AuthDepends` → `DataScopeDepends`);gRPC 用 server interceptor 注入身份元数据
- **应用服务层**:编排查询 / 计算掌握度 / 评估预警 / 发布事件,不直接访问存储
- **数据访问层**:ClickHouse 查询封装(含 `FINAL` / `argMax` 去重)+ Kafka producer + gRPC client(调 iam)+ Redis client
- **CDC 消费者**:独立后台任务(lifespan 启动),与 HTTP/gRPC 入口解耦;P6 演进为多实例水平扩展(ExamCache Redis 化)
## 2. 领域模型
data-ana 是**纯读模型服务**,不持有写聚合根。领域模型为**视图聚合**(ClickHouse 物化):
### 聚合根(视图型)
| 聚合根 | 含义 | 物化载体 | 不变式 |
| ------------------ | ---------------- | ------------------------------------------- | ------------------------------------------------------------- |
| `StudentDashboard` | 学生学情宽表 | ClickHouse `student_dashboard_view` | 同一 (student_id, exam_id, knowledge_point_id) 仅保留最新版本 |
| `ClassPerformance` | 班级成绩聚合 | ClickHouse 即时聚合(不物化) | 聚合维度为 class_id + 时间窗 |
| `StudentErrorBook` | 学生错题本 | ClickHouse `student_errors` | 同一 (student_id, question_id) 累计 error_count |
| `MasterySnapshot` | 知识点掌握度快照 | ClickHouse `mastery_snapshot`(**新增**) | 同一 (student_id, knowledge_point_id) 保留历史版本 |
| `AttendanceLog` | 学生考勤记录 | ClickHouse `attendance_logs`(**v2 新增**) | 同一 (student_id, class_id, attendance_date) 仅保留最新版本 |
| `AiUsageLog` | AI 用量计费记录 | ClickHouse `ai_usage_log`(**v2 新增**) | 同一 request_id 唯一(消费方按 request_id 幂等) |
> **ai-allocation §5 对齐**:5 张宽表完整覆盖(student_dashboard_view / student_errors / mastery_snapshot / attendance_logs / ai_usage_log)。exam_results / homework_submissions / grade_records 数据已折叠进 student_dashboard_view 宽表(按 event_type 区分),避免宽表数量爆炸。
### 值对象
- `WeakPoint`:knowledge_point_id + title + mastery_level
- `TrendPoint`:date + score
- `Warning`:warning_type (LOW_MASTERY / SCORE_DROP / ABSENT_FREQUENT) + target_id + threshold + current_value + severity (INFO / WARN / CRITICAL)
- `DataScope`:level (SELF/CLASS/GRADE/SCHOOL/DISTRICT/ALL) + scope_ids(具体可见的 class_id / grade_id 列表)
- `MasteryMethod`:枚举(WEIGHTED_MOVING_AVG / SIMPLE_AVG / FORGETTING_CURVE),见 §9
### 聚合间通信
- 同服务内:直接函数调用(Application Service → Repository)
- 跨服务:仅通过 Kafka 事件(发布 `mastery.updated`)+ gRPC(调 iam 查 DataScope;**待 P5+** ai 反向调 data-ana 查学情)
## 3. 数据模型(ClickHouse DDL)
> DDL 文件由 coord 统一管理在 `infra/clickhouse/ddl/`(待 coord 建立),data-ana 提供内容。
### 3.1 宽表 `student_dashboard_view`
```sql
-- 学生学情宽表:每次成绩写入产生一行,按 ORDER BY 去重保留最新版本
CREATE TABLE IF NOT EXISTS student_dashboard_view
(
student_id String,
class_id String,
exam_id String,
subject_id String,
score Float64,
rank_in_class UInt32,
knowledge_point_id String,
mastery_level Float32, -- 0.0-1.0
error_count UInt32,
last_updated DateTime64(3, 'UTC')
)
ENGINE = ReplacingMergeTree(last_updated)
PARTITION BY toYYYYMM(last_updated)
ORDER BY (student_id, exam_id, knowledge_point_id)
SETTINGS index_granularity = 8192;
```
**索引策略**:
- ORDER BY `(student_id, exam_id, knowledge_point_id)`:主键索引,支持按学生查学情、按考试查成绩、按知识点查掌握度
- PARTITION BY `toYYYYMM(last_updated)`:按月分区,支持历史数据归档
- ReplacingMergeTree(last_updated):同 ORDER BY 自动去重,保留 last_updated 最大版本(幂等消费保证)
### 3.2 错题本 `student_errors`
```sql
CREATE TABLE IF NOT EXISTS student_errors
(
student_id String,
question_id String,
knowledge_point_id String,
error_count UInt32,
last_error_time DateTime64(3, 'UTC'),
content String
)
ENGINE = ReplacingMergeTree(last_error_time)
PARTITION BY toYYYYMM(last_error_time)
ORDER BY (student_id, question_id);
```
### 3.3 掌握度快照 `mastery_snapshot`(**新增**)
```sql
-- 知识点掌握度历史快照:每次掌握度计算产生新版本,支持趋势查询
CREATE TABLE IF NOT EXISTS mastery_snapshot
(
student_id String,
knowledge_point_id String,
subject_id String,
mastery_level Float32,
calculated_at DateTime64(3, 'UTC'),
calculation_method LowCardinality(String) -- 'weighted_moving_avg' / 'simple_avg'
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(calculated_at)
ORDER BY (student_id, knowledge_point_id, calculated_at);
```
### 3.4 用量计费 `ai_usage_log`(**新增,供 ai 服务写入**)
```sql
-- AI 用量记录:ai 服务通过 Kafka 事件投递,data-ana 消费落库
CREATE TABLE IF NOT EXISTS ai_usage_log
(
request_id String,
user_id String,
provider LowCardinality(String), -- 'openai' / 'anthropic' / 'baichuan' / 'local'
model LowCardinality(String),
prompt_tokens UInt32,
completion_tokens UInt32,
total_tokens UInt32,
latency_ms UInt32,
success Boolean,
cost_cents UInt32, -- 计费(分),便于聚合
occurred_at DateTime64(3, 'UTC')
)
ENGINE = ReplacingMergeTree(occurred_at)
PARTITION BY toYYYYMM(occurred_at)
ORDER BY (request_id); -- 按 request_id 幂等去重
```
### 3.5 考勤记录 `attendance_logs`(**v2 新增,ai-allocation §5 第 5 张宽表**)
```sql
-- 学生考勤记录:core-edu attendance 表 CDC 同步
CREATE TABLE IF NOT EXISTS attendance_logs
(
student_id String,
class_id String,
attendance_date Date,
status LowCardinality(String), -- 'present' / 'absent' / 'late' / 'leave'
recorded_by String, -- 教师用户 ID
remark String DEFAULT '',
occurred_at DateTime64(3, 'UTC')
)
ENGINE = ReplacingMergeTree(occurred_at)
PARTITION BY toYYYYMM(attendance_date)
ORDER BY (student_id, class_id, attendance_date);
```
**索引策略**:
- ORDER BY `(student_id, class_id, attendance_date)`:支持按学生查考勤历史、按班级查当日考勤、按日期范围统计
- PARTITION BY `toYYYYMM(attendance_date)`:按月分区,支持历史归档
- ReplacingMergeTree(occurred_at):同一记录多版本去重(教师修改考勤状态时产生新版本)
### 3.6 读写分离策略
| 操作 | 路径 | 说明 |
| -------------- | ----------------------------------- | ------------------------------------------ |
| 学情查询 | ClickHouse 宽表 | 实时聚合,亚秒级响应 |
| 错题本查询 | ClickHouse student_errors | 实时查询 |
| 掌握度趋势 | ClickHouse mastery_snapshot | 历史快照 |
| 考勤查询 | ClickHouse attendance_logs | 学生考勤历史 |
| AI 用量统计 | ClickHouse ai_usage_log | 按用户/模型/时间窗聚合 |
| 掌握度计算 | CDC 触发 → 内存计算 → 写 ClickHouse | 派生数据,非事务写 |
| DataScope 解析 | gRPC 调 iam | 实时查询,结果 Redis 缓存 5min(004 §6.3) |
> **双轨读边界**(004 §1.1a):BFF 实时查主库走 core-edu,data-ana 仅承担聚合查宽表职责,二者不重叠。
## 4. API 设计
### 4.1 HTTP 端点(保留作 Gateway 直连降级)
| method | path | 权限 | 请求 | 响应 | 说明 |
| ------ | -------------------------------------------- | ----------------------------- | ------------------------------------------------------ | ---------------------------------------------------------------------- | ----------------------------------- |
| GET | `/healthz` | — | — | `{status, service}` | liveness |
| GET | `/readyz` | — | — | `{status, ready, degraded, clickhouse, cdc_consumer, redis, iam_grpc}` | readiness(v2 补 redis / iam_grpc) |
| GET | `/metrics` | — | — | Prometheus 格式 | 指标 |
| GET | `/analytics/class/{class_id}/performance` | `ANALYTICS_CLASS_READ` | query: `subject_id?`, `start_date?`, `end_date?` | `ActionState` | 班级成绩分析 |
| GET | `/analytics/student/{student_id}/weakness` | `ANALYTICS_STUDENT_READ` | query: `subject_id?` | `ActionState` | 学生薄弱知识点(DataScope 过滤) |
| GET | `/analytics/student/{student_id}/errorbook` | `ANALYTICS_STUDENT_READ` | query: `page?`, `size?` | `ActionState` | 学生错题本(DataScope 过滤) |
| GET | `/analytics/student/{student_id}/trend` | `ANALYTICS_STUDENT_READ` | query: `start_date`, `end_date`, `subject_id?` | `ActionState` | 学习趋势 |
| GET | `/analytics/student/{student_id}/attendance` | `ANALYTICS_STUDENT_READ` | query: `start_date`, `end_date` | `ActionState` | 学生考勤历史(**v2 新增**) |
| GET | `/analytics/dashboard/teacher/{user_id}` | `ANALYTICS_TEACHER_DASHBOARD` | query: `class_id?` | `ActionState` | 教师仪表盘聚合 |
| GET | `/analytics/dashboard/student/{user_id}` | `ANALYTICS_STUDENT_DASHBOARD` | — | `ActionState` | 学生仪表盘(**v2 新增**) |
| GET | `/analytics/dashboard/parent/{user_id}` | `ANALYTICS_PARENT_DASHBOARD` | query: `student_id` | `ActionState` | 家长仪表盘(**v2 新增**) |
| GET | `/analytics/dashboard/admin/{user_id}` | `ANALYTICS_ADMIN_DASHBOARD` | query: `scope?`, `scope_id?` | `ActionState` | 管理员仪表盘(**v2 新增**) |
| GET | `/analytics/warnings` | `ANALYTICS_WARNING_READ` | query: `class_id?`, `severity?`, `since?` | `ActionState` | 预警列表(**v2 新增**) |
| GET | `/analytics/mastery/distribution` | `ANALYTICS_CLASS_READ` | query: `class_id`, `subject_id`, `knowledge_point_id?` | `ActionState` | 班级掌握度分布(**v2 新增**) |
### 4.2 gRPC 契约(analytics.proto,**v2 扩展**)
| RPC | 请求 | 响应 | 权限 |
| ------------------------ | -------------------------------------------------------------------------- | --------------------------- | ----------------------------- |
| `GetClassPerformance` | `GetClassPerformanceRequest{class_id, subject_id, start_date, end_date}` | `ClassPerformance` | `ANALYTICS_CLASS_READ` |
| `GetStudentWeakness` | `GetStudentWeaknessRequest{student_id, subject_id}` | `StudentWeakness` | `ANALYTICS_STUDENT_READ` |
| `GetLearningTrend` | `GetLearningTrendRequest{student_id, start_date, end_date, subject_id?}` | `LearningTrend` | `ANALYTICS_STUDENT_READ` |
| `GetTeacherDashboard` | `GetTeacherDashboardRequest{user_id, class_id?}` | `TeacherDashboard` | `ANALYTICS_TEACHER_DASHBOARD` |
| `GetStudentDashboard` | `GetStudentDashboardRequest{user_id}` | `StudentDashboard` | `ANALYTICS_STUDENT_DASHBOARD` |
| `GetParentDashboard` | `GetParentDashboardRequest{user_id, student_id}` | `ParentDashboard` | `ANALYTICS_PARENT_DASHBOARD` |
| `GetAdminDashboard` | `GetAdminDashboardRequest{user_id, scope, scope_id?}` | `AdminDashboard` | `ANALYTICS_ADMIN_DASHBOARD` |
| `GetWarnings` | `GetWarningsRequest{class_id?, severity?, since?}` | `WarningList` | `ANALYTICS_WARNING_READ` |
| `GetMasteryDistribution` | `GetMasteryDistributionRequest{class_id, subject_id, knowledge_point_id?}` | `MasteryDistribution` | `ANALYTICS_CLASS_READ` |
| `SubscribeMasteryUpdate` | `SubscribeMasteryUpdateRequest{student_id?, class_id?}` | `stream MasteryUpdateEvent` | `ANALYTICS_STUDENT_READ` |
> **proto 扩展提案**(待 coord 审议):上述新增 RPC 的 message 定义需在 `packages/shared-proto/proto/analytics.proto` 补充。`SubscribeMasteryUpdate` 为 server-streaming RPC,为 P5+ AI 个性化推荐预留实时推送通道。
**权限校验**:gRPC server interceptor 从 metadata 提取 `x-user-id` / `x-user-roles` / `x-data-scope`,调用 `AuthDepends` 等价逻辑。
### 4.3 ActionState 统一响应信封(004 §11.5 P0 整改)
所有 HTTP/gRPC 响应必须遵循 ActionState 信封:
```python
from pydantic import BaseModel
from typing import Generic, TypeVar
T = TypeVar("T")
class ActionStateError(BaseModel):
code: str
message: str
details: dict[str, Any] | None = None
trace_id: str | None = None
class ActionState(BaseModel, Generic[T]):
"""统一响应信封(004 §11.5).
成功:{success: true, data: T}
失败:{success: false, error: {code, message, details?, trace_id?}}
降级:success=true 但 error.details.degraded=true(保留功能但数据可能不完整)
"""
success: bool
data: T | None = None
error: ActionStateError | None = None
@classmethod
def ok(cls, data: T, *, degraded: bool = False) -> "ActionState[T]":
details = {"degraded": True} if degraded else None
# 简化:degraded 作为 data 的元信息附加,复杂场景用 details
return cls(success=True, data=data, error=None if not degraded else
ActionStateError(code="DATA_ANA_DEGRADED", message="degraded mode",
details=details))
@classmethod
def fail(cls, code: str, message: str, *, trace_id: str | None = None,
details: dict[str, Any] | None = None) -> "ActionState[T]":
return cls(success=False, data=None,
error=ActionStateError(code=code, message=message,
details=details, trace_id=trace_id))
```
```python
# 示例:班级成绩分析响应
class ClassPerformanceData(BaseModel):
class_id: str
average_score: float
pass_rate: float
total_students: int
scores: list[StudentScore] = [] # 详细成绩列表(受 DataScope 过滤)
class StudentScore(BaseModel):
student_id: str
score: float
grade: str
# 端点返回类型:ActionState[ClassPerformanceData]
```
> **P0 整改要点**:当前 main.py 返回 `{success, data, degraded}` 三字段平铺,违反 004 §11.5。实现阶段需重构为 `ActionState[T]` 泛型,`degraded` 移到 `error.details.degraded`。
## 5. 事件设计
### 5.1 消费的事件
| Topic | 来源 | 消息格式 | 消费动作 |
| ------------------------------------------------------ | ----------------------------------------- | --------------------------------------------- | ---------------------------------------------------------------------------- |
| `edu-cdc.next_edu_cloud.core_edu_grades` | Debezium CDC(core-edu MySQL) | Debezium JSON(before/after/source/op/ts_ms) | 解析 → 查 ExamCache 填 class_id → 计算掌握度 → upsert student_dashboard_view |
| `edu-cdc.next_edu_cloud.core_edu_exams` | Debezium CDC | Debezium JSON | upsert ExamCache(exam_id → class_id, subject_id) |
| `edu-cdc.next_edu_cloud.core_edu_homework_submissions` | Debezium CDC | Debezium JSON | 记录作业提交行为 → 更新 student_dashboard_view |
| `edu-cdc.next_edu_cloud.core_edu_attendance` | Debezium CDC(**v2 新增**) | Debezium JSON | 落 attendance_logs 表 |
| `edu-cdc.next_edu_cloud.classes` | Debezium CDC | Debezium JSON | 同步班级维度(head_teacher_id)用于教师 DataScope |
| `edu-cdc.next_edu_cloud.iam_users` | Debezium CDC | Debezium JSON | 同步用户 dataScope 用于查询过滤(避免每次查 iam) |
| `edu-cdc.next_edu_cloud.content_knowledge_points` | Debezium CDC(**v2 新增**,content 服务) | Debezium JSON | 同步知识点标题/前置关系到 mastery_snapshot 元数据 |
| `edu.insight.ai.usage` | ai 服务 Kafka producer | JSON(UsageRecord) | 落 ai_usage_log 表 |
> **004 §4.1 对齐**:content → data-ana 事件流已声明("教学内容变更通知"),通过 CDC 落知识点元数据。当 content 服务发布 `KnowledgePointUpdated` 领域事件时(待 P4 后期评估),可切换为订阅 `edu.content.events` topic,避免依赖 content MySQL binlog。
**幂等性**:
- CDC 事件:依赖 ClickHouse `ReplacingMergeTree(last_updated)` 引擎按 ORDER BY 去重 + ReplacingMergeTree 的 version 列
- 领域事件(若消费):基于 `event_id` 去重(Redis SETNX `data_ana:dedup:{event_id}` TTL 7 天)
- 消费者 offset:v1 用 `enable_auto_commit=True` 简化;**v2 改进**:手动 commit(at-least-once),commit 前确保 ClickHouse 写入成功
### 5.2 发布的事件
| 事件 | Topic | 触发时机 | 消费者 | Payload |
| ------------------ | ---------------------------------------------- | ------------------------------------------------------------------ | --------------------------------------------- | -------------------------------------------------------------------------------------- |
| `MasteryUpdated` | `edu.insight.mastery.updated` | 掌握度计算完成(CDC grades 事件触发后异步计算) | core-edu(推荐个性化练习)、msg(掌握度预警) | `{event_id, student_id, knowledge_point_id, mastery_level, calculated_at}` |
| `WarningTriggered` | `edu.insight.warning.triggered`(**v2 新增**) | 预警阈值触发(掌握度低于 0.4 / 成绩环比下降 20% / 缺勤 ≥ 3 次/周) | msg(推送通知)、core-edu(标记关注) | `{event_id, warning_type, target_id, threshold, current_value, severity, occurred_at}` |
**发布实现**(004 §12.2 + §15.3 #6 已仲裁:派生数据事件豁免 Outbox):
- 掌握度计算、预警评估均为**派生数据**(非业务事务写),不需要 Outbox 保证事务一致
- 直接用 `aiokafka.AIOKafkaProducer` 发布,`idempotent=true` + `transactional_id="data-ana-producer"`
- 失败重试 3 次,仍失败记录日志 + 落 `publish_failed` 本地表(SQLite 文件,P6 评估是否引入 Outbox)
- 事件 schema 版本化:`schema_version` 字段(v1/v2/...),禁止破坏性变更,见 §17
### 5.3 CDC 水平扩展策略(v2 新增,P6 铺垫)
**P4 单实例模式**:
- 1 个 CdcConsumer 实例消费所有 topic,consumer group = `data-ana-cdc`
- ExamCache 内存 LRU(max 10000 条)
- 风险:单点故障 + 重启 lag 累积
**P6 多实例水平扩展**:
- consumer group 不变,多实例分摊 partition
- ExamCache → Redis 化(key: `data_ana:exam:{exam_id}` TTL 30 天)
- 幂等去重 → Redis SETNX(已具备,多实例天然共享)
- partition 数从 P4 的 1 → P6 的 N(按 topic 流量评估)
- 监控:consumer lag 指标(已具备 `data_ana_cdc_consumer_lag`)+ HPA 阈值
**演进路径**:
```
P4: 1 实例 + 内存 ExamCache + auto commit
↓
P5: 1 实例 + Redis ExamCache + 手动 commit(过渡)
↓
P6: N 实例 + Redis ExamCache + 手动 commit + HPA
```
> 单实例 → 多实例的代码改动量:仅 ExamCache 实现(dict → Redis client),其余抽象层不变。
## 6. 横切关注点对齐清单
### 6.1 权限装饰器等价物(FastAPI Depends)
```python
# 权限点常量(与 iam 权限点对齐,v2 扩展 4 端 Dashboard + 预警权限)
class Permissions:
ANALYTICS_CLASS_READ = "analytics:class:read"
ANALYTICS_STUDENT_READ = "analytics:student:read"
ANALYTICS_TEACHER_DASHBOARD = "analytics:teacher:dashboard"
ANALYTICS_STUDENT_DASHBOARD = "analytics:student:dashboard" # v2 新增
ANALYTICS_PARENT_DASHBOARD = "analytics:parent:dashboard" # v2 新增
ANALYTICS_ADMIN_DASHBOARD = "analytics:admin:dashboard" # v2 新增
ANALYTICS_WARNING_READ = "analytics:warning:read" # v2 新增
async def require_permission(permission: str) -> UserContext:
"""FastAPI Depends 权限校验.
从 x-user-id / x-user-roles 头提取身份,校验角色是否含 permission。
"""
...
async def inject_data_scope(ctx: UserContext = Depends(require_permission(...)))-> DataScope:
"""注入 DataScope(从 iam.getEffectiveDataScope 查询,Redis 缓存 5min)."""
...
```
### 6.2 错误码清单(前缀 `DATA_ANA_*`,v2 扩展)
| 错误码 | 触发条件 | HTTP | gRPC status |
| --------------------------------- | --------------------------------------- | ------------------- | ------------------ |
| `DATA_ANA_UNAUTHORIZED` | 缺失 x-user-id 头或 token 无效 | 401 | UNAUTHENTICATED |
| `DATA_ANA_FORBIDDEN` | 角色无对应权限 | 403 | PERMISSION_DENIED |
| `DATA_ANA_DATASCOPE_VIOLATION` | 查询目标超出 DataScope 范围 | 403 | PERMISSION_DENIED |
| `DATA_ANA_CLICKHOUSE_UNAVAILABLE` | ClickHouse 不可达(降级模式仍返回骨架) | 200 + degraded:true | OK + degraded flag |
| `DATA_ANA_KAFKA_UNAVAILABLE` | Kafka producer 不可达(事件发布失败) | 200 + degraded | OK + degraded |
| `DATA_ANA_REDIS_UNAVAILABLE` | Redis 不可达(降级直查 iam) | 200 + degraded | OK + degraded |
| `DATA_ANA_IAM_GRPC_UNAVAILABLE` | iam gRPC 不可达(fallback SELF 范围) | 200 + degraded | OK + degraded |
| `DATA_ANA_INVALID_DATE_RANGE` | start_date > end_date | 400 | INVALID_ARGUMENT |
| `DATA_ANA_STUDENT_NOT_FOUND` | student_id 不存在 | 404 | NOT_FOUND |
| `DATA_ANA_CLASS_NOT_FOUND` | class_id 不存在 | 404 | NOT_FOUND |
| `DATA_ANA_WARNING_NOT_FOUND` | 预警 ID 不存在 | 404 | NOT_FOUND |
| `DATA_ANA_RATE_LIMITED` | 触发限流(004 §4.3 限流策略) | 429 | RESOURCE_EXHAUSTED |
| `DATA_ANA_INTERNAL_ERROR` | 未捕获异常 | 500 | INTERNAL |
### 6.3 Logger 初始化
- 位置:`main.py` `init_logger()`(已具备)
- 配置:`structlog.make_filtering_bound_logger(level)` + `TimeStamper(fmt="iso")` + `ConsoleRenderer`
- **改进**:生产环境改用 `structlog.processors.JSONRenderer()`(当前 ConsoleRenderer 适合开发)
- 日志字段:`request_id` / `user_id` / `trace_id` / `service="data-ana"`(004 §10 强制)
### 6.4 Metrics 指标清单(v2 扩展业务指标)
| 指标名 | 类型 | 标签 | 描述 |
| --------------------------------------------- | --------- | ---------------------- | ---------------------------- |
| `data_ana_http_requests_total` | Counter | method, path, status | HTTP 请求总数 |
| `data_ana_http_request_duration_seconds` | Histogram | method, path | HTTP 请求延迟 |
| `data_ana_grpc_requests_total` | Counter | rpc, status | gRPC 请求总数(v2 新增) |
| `data_ana_grpc_request_duration_seconds` | Histogram | rpc | gRPC 请求延迟(v2 新增) |
| `data_ana_clickhouse_query_duration_seconds` | Histogram | query_type | ClickHouse 查询延迟 |
| `data_ana_clickhouse_query_total` | Counter | query_type, status | ClickHouse 查询总数 |
| `data_ana_cdc_events_consumed_total` | Counter | table, op | CDC 事件消费总数 |
| `data_ana_cdc_event_process_duration_seconds` | Histogram | table | CDC 事件处理延迟 |
| `data_ana_cdc_consumer_lag` | Gauge | topic, partition | CDC 消费者 lag |
| `data_ana_mastery_calculated_total` | Counter | method | 掌握度计算次数(按算法) |
| `data_ana_mastery_published_total` | Counter | status | mastery.updated 事件发布数 |
| `data_ana_warning_triggered_total` | Counter | warning_type, severity | 预警触发数(v2 新增) |
| `data_ana_datascope_cache_hits_total` | Counter | — | DataScope 缓存命中 |
| `data_ana_redis_ops_total` | Counter | op, status | Redis 操作数(v2 新增) |
| `data_ana_iam_grpc_calls_total` | Counter | status | iam gRPC 调用数(v2 新增) |
| `data_ana_ai_usage_ingested_total` | Counter | provider, model | AI 用量事件消费数(v2 新增) |
### 6.5 Tracer 初始化
- 位置:`main.py` `init_tracer()`(已具备)
- endpoint:`settings.otel_endpoint` + `/v1/traces`
- **改进**:gRPC server 注册 `grpc.aio.ServerInterceptor` 透传 W3C trace context
- **改进**:CDC 消费者 span 注入 trace context(从 Kafka header 提取 traceparent)
### 6.6 /healthz 检查逻辑
- liveness:仅进程存活(已具备)
### 6.7 /readyz 检查逻辑(v2 扩展)
```python
async def readyz() -> dict:
return {
"status": "ok" if all_ready else "degraded",
"ready": all_ready,
"degraded": not all_ready,
"clickhouse": "ok" | "unreachable" | "not_configured",
"cdc_consumer": "running" | "disabled" | "failed",
"kafka_brokers": settings.kafka_brokers or None,
"iam_grpc": "ok" | "unreachable", # v2 新增
"redis": "ok" | "unreachable", # v2 新增
"timestamp": datetime.now(UTC).isoformat(),
}
```
### 6.8 优雅关闭顺序
1. HTTP server stop accepting new requests(uvicorn graceful shutdown)
2. gRPC server graceful stop(等待在途 RPC 完成,30s 超时)
3. CDC consumer stop(等待在途消息处理完成,commit offset)
4. Kafka producer flush + close(确保 mastery.updated 事件已投递)
5. ClickHouse client close
6. iam gRPC channel close
**信号处理**:注册 `signal.SIGTERM` handler,触发上述顺序。
## 7. 与其他模块的交互点(契约清单,v2 扩展)
| 方向 | 对方服务 | 协议 | 接口/事件 | 用途 |
| ------ | -------------------------------------- | ----- | --------------------------------------------------------------------------------- | ------------------------------------ |
| 被调用 | api-gateway | HTTP | `/analytics/*` | Gateway 代理(降级通道) |
| 被调用 | teacher-bff / student-bff / parent-bff | gRPC | `AnalyticsService.*` | BFF 聚合查询(4 端 Dashboard) |
| 被调用 | ai(P5+) | gRPC | `AnalyticsService.GetStudentWeakness / GetLearningTrend / SubscribeMasteryUpdate` | AI 个性化出题上下文 + 实时掌握度推送 |
| 调用 | iam | gRPC | `IamService.GetEffectiveDataScope`(004 §15.3 #5 已仲裁 P4 补全) | DataScope 解析 |
| 消费 | core-edu(CDC) | Kafka | `edu-cdc.next_edu_cloud.core_edu_grades/exams/homework_submissions/attendance` | 学情 + 考勤数据投递 |
| 消费 | core-edu(CDC) | Kafka | `edu-cdc.next_edu_cloud.classes` | 班级维度同步 |
| 消费 | iam(CDC) | Kafka | `edu-cdc.next_edu_cloud.iam_users` | 用户 dataScope 同步 |
| 消费 | content(CDC) | Kafka | `edu-cdc.next_edu_cloud.content_knowledge_points` | 知识点元数据同步(v2 新增) |
| 消费 | ai(P5+) | Kafka | `edu.insight.ai.usage`(004 §15.3 #6 已仲裁,004 §7.2 已登记) | AI 用量落库 |
| 发布 | core-edu / msg | Kafka | `edu.insight.mastery.updated`(004 §12.2 + §15.3 #6 已仲裁:派生数据豁免 Outbox) | 掌握度更新通知 |
| 发布 | msg / core-edu | Kafka | `edu.insight.warning.triggered`(v2 新增,需 coord 在 004 §7.2 登记) | 预警触发通知 |
> **004 §4.1 服务间通信矩阵对齐**:content → data-ana 已声明"教学内容变更通知",本次落实为 CDC 订阅 `content_knowledge_points` 表。
## 8. 风险与假设
### 8.1 假设
- **假设 1**:iam 在 P4 阶段提供 `GetEffectiveDataScope(userId) → DataScope` gRPC API(004 §15.3 #5 已仲裁)。若 iam 未及时提供,fallback 为:从 `x-user-roles` 头推导(admin=ALL, teacher=CLASS_TAUGHT, student=SELF),但无法支持细粒度年级/学校范围
- **假设 2**:core-edu 的 `core_edu_homework_submissions` / `core_edu_attendance` 表存在 binlog。若不存在,需 core-edu 补表或走 Outbox 事件(领域事件备通道)
- **假设 3**:ClickHouse `ReplacingMergeTree` 在查询时需 `FINAL` 关键字确保去重生效。**v2 修复要求**:所有查询加 `FINAL` 或使用 `argMax` 聚合(当前实现未加,是 P0 整改项)
- **假设 4**:coord 已在 004 §7.2 登记新增 `edu.insight.ai.usage` topic(004 §15.3 #6 已仲裁)
### 8.2 技术风险
| 风险 | 影响 | 缓解措施 |
| --------------------------- | ------------------------------------ | ------------------------------------------------------------ |
| ClickHouse 查询延迟超 5s | 违反 P4 退出标准 | 宽表索引优化 + 物化视图预聚合 + 查询超时 3s 降级 |
| CDC 消费者 lag 过大 | 学情数据延迟 > 5s | 监控 lag + 告警 + 水平扩展消费者(分区数提升) |
| ExamCache 内存泄漏 | 长期运行 OOM | LRU 淘汰策略(max 10000 条)+ 定期清理过期 exam;P6 迁 Redis |
| mastery.updated 事件丢失 | 下游 core-edu/msg 收不到通知 | Kafka producer `acks=all` + 本地失败表重试 |
| iam gRPC 不可达 | DataScope 无法解析 → 查询降级为 SELF | Redis 缓存 5min + fallback SELF 范围(最保守) |
| ClickHouse `FINAL` 查询性能 | 查询变慢 | 使用 `argMax` 替代 `FINAL`,或在写入时去重 |
| ClickHouse 单点故障 | 全部查询不可用 | P6 引入 ClickHouse 集群(ReplicatedMergeTree)+ 副本 |
| 多实例 CDC 重复消费 | ClickHouse 重复写入 | ReplacingMergeTree ORDER BY 去重 + Redis SETNX 幂等 |
| 大数据量 Dashboard 聚合超时 | 管理员 Dashboard 全校聚合慢 | 物化视图预聚合 + 异步刷新 + 缓存 5min |
| 知识点元数据与成绩关联失败 | mastery_snapshot 缺知识点标题 | content CDC 同步 + 缺失时显示 `knowledge_point_id` |
### 8.3 coord 交叉审查结论对齐(v2:原"未决设计决策"已全部仲裁)
| # | 议题 | coord 仲裁结论 | 涉及文档 |
| --- | ------------------------------------------------------------------------------------ | ----------------------------------------------------------------- | ----------------------- |
| 1 | data-ana 发布 `edu.insight.mastery.updated` 用直接 producer(非 Outbox)是否合规 | ✅ 已仲裁:派生数据事件豁免 Outbox(004 §12.2 + §15.3 #6) | 004 §12.2 |
| 2 | 新增 `edu.insight.ai.usage` topic + `AIUsageEvent` proto message | ✅ 已仲裁:coord 在 004 §7.2 登记 + events.proto 补 message | 004 §7.2 + events.proto |
| 3 | iam 新增 `GetEffectiveDataScope` gRPC RPC | ✅ 已仲裁:P4 补全(004 §15.3 #5) | iam.proto |
| 4 | data-ana 实现 gRPC server 决策 | ✅ 已仲裁:P4 启用(004 §4.2 gRPC 启用阶段矩阵) | 004 §4.2 |
| 5 | ClickHouse DDL 管理位置(建议 `infra/clickhouse/ddl/`) | ⏳ 待 coord 建立 `infra/clickhouse/ddl/` 目录 + data-ana 提供内容 | infra/ |
| 6 | 端口冲突检查:data-ana HTTP=3006 / gRPC=50055 | ✅ 无冲突(004 §1.2 端口分配) | 004 §1.2 |
| 7 | 错误码前缀检查:`DATA_ANA_*` | ✅ 无冲突(004 §11.4 错误码前缀矩阵) | 004 §11.4 |
| 8 | 黄金模板对齐:Python 服务无 NestJS 装饰器,权限校验用 FastAPI Depends 等价物是否认可 | ✅ 已仲裁:认可(004 §15.3) | — |
| 9 | `edu.insight.warning.triggered` topic 新增 | ⏳ 待 coord 在 004 §7.2 登记(v2 新提案) | 004 §7.2 |
| 10 | analytics.proto 扩展(4 端 Dashboard / Warning / MasteryDistribution / Stream RPC) | ⏳ 待 coord 审议(v2 新提案,§4.2) | analytics.proto |
> 所有"未决设计决策"已消除,进入实现阶段无阻塞。剩余 3 项 ⏳ 为 v2 新提案,待 coord 审议但不阻塞 P4 主体实现(可先用现有 RPC + HTTP 端点兜底)。
---
## 9. 掌握度计算算法(ai-allocation §5 + pending-features P4 要求)
### 9.1 算法选型
| 算法 | 公式 | 适用场景 | 默认 |
| --------------------- | ------------------------------------------------------------ | ---------------------------- | ---------------- |
| `WEIGHTED_MOVING_AVG` | $M = \sum_{i=1}^{n} w_i \cdot s_i / \sum w_i$,$w_i = 0.6^i$ | 学生有多次同知识点考试成绩 | ✅ |
| `FORGETTING_CURVE` | $M(t) = M_0 \cdot e^{-t/S}$,$S$ = 30 天 | 长期未考的知识点,掌握度衰减 | ✅(与上者叠加) |
| `SIMPLE_AVG` | $M = \frac{1}{n}\sum s_i$ | 数据稀疏(< 3 次) | 兜底 |
**综合公式**(默认策略 `WEIGHTED_MOVING_AVG + FORGETTING_CURVE`):
```
M_final = max(
WEIGHTED_MOVING_AVG(近 5 次成绩, 时间倒序加权),
FORGETTING_CURVE(最近一次掌握度, 距今天数)
)
```
- 取 max 是为了避免"刚考完高分但长期没考"被遗忘曲线过度拉低
- 近 5 次不足时,用 SIMPLE_AVG 兜底
- 所有成绩归一化到 [0, 1](score / 100)
### 9.2 算法参数(可配置)
| 参数 | 默认值 | 说明 |
| ---------------------- | ------ | ---------------------------- |
| `window_size` | 5 | 加权滑动窗口大小 |
| `decay_base` | 0.6 | 加权衰减基数(越近权重越高) |
| `forgetting_half_life` | 30 | 遗忘曲线半衰期(天) |
| `min_samples_for_wma` | 3 | 低于此数降级为 SIMPLE_AVG |
| `score_max` | 100 | 成绩满分(用于归一化) |
### 9.3 触发时机
1. CDC `core_edu_grades` 事件到达 → 异步计算 → 写 `mastery_snapshot` + `student_dashboard_view` → 发布 `MasteryUpdated`
2. 定时任务(每日 02:00)→ 对所有知识点重算遗忘曲线衰减 → 更新 `mastery_snapshot`(不发布事件,避免风暴)
3. 手动触发(管理员 Dashboard "重算掌握度"按钮,v2 预留 API `POST /analytics/mastery/recalculate`)
### 9.4 预警阈值设计(ai-allocation §5 要求)
| 预警类型 | 触发条件 | severity | 通知动作 |
| ----------------- | ---------------------- | -------- | ---------------------- |
| `LOW_MASTERY` | mastery_level < 0.4 | WARN | 推送学生 + 教师 |
| `CRITICAL_LOW` | mastery_level < 0.2 | CRITICAL | 推送学生 + 教师 + 家长 |
| `SCORE_DROP` | 本次成绩环比下降 ≥ 20% | WARN | 推送教师 |
| `ABSENT_FREQUENT` | 单周缺勤 ≥ 3 次 | WARN | 推送教师 + 家长 |
| `TREND_DECLINE` | 连续 3 次成绩下降 | INFO | 教师仪表盘标记 |
**预警去重**:同一 (target_id, warning_type) 24h 内只触发一次(Redis 位图 `data_ana:warning:{target_id}:{type}:{date}`)
## 10. 缓存策略(004 §6.3 缓存策略矩阵对齐)
### 10.1 Redis 连接
```python
# config.py 新增
redis_url: str = "redis://edu-redis:6379/3" # data-ana 专用 db=3
redis_pool_size: int = 10
redis_socket_timeout_ms: int = 200
```
### 10.2 缓存键命名规范
| 键模式 | TTL | 用途 | 失效策略 |
| -------------------------------------------- | ----- | -------------------------- | ---------------------- |
| `data_ana:datascope:{user_id}` | 5min | DataScope 缓存 | 角色变更主动 DEL |
| `data_ana:exam:{exam_id}` | 30day | ExamCache(P6 多实例共享) | exam 结束 30 天后过期 |
| `data_ana:dedup:{event_id}` | 7day | 事件幂等去重 | 自然过期 |
| `data_ana:warning:{target_id}:{type}:{date}` | 25h | 预警去重位图 | 自然过期 |
| `data_ana:dashboard:{user_id}:{scope}` | 5min | Dashboard 聚合结果缓存 | 主动 DEL on 新成绩写入 |
| `data_ana:kp_meta:{knowledge_point_id}` | 1day | 知识点元数据缓存 | content CDC 更新时 DEL |
### 10.3 缓存降级
- Redis 不可达:所有缓存查询跳过(cache miss),直查 ClickHouse / iam gRPC
- DataScope 缓存 miss → 调 iam gRPC(增加 iam 负载,需监控)
- Dashboard 缓存 miss → 直查 ClickHouse(增加 CH 负载,可能触发查询超时降级)
## 11. CDC 水平扩展(P6 铺垫,详见 §5.3)
P4 单实例 → P6 多实例演进路径已在 §5.3 详述。本节补充容量规划。
### 11.1 ClickHouse 容量规划
**假设**:全校 5000 学生,每学生每月 20 次成绩写入,5 个知识点/次。
| 宽表 | 月增量(行) | 单行大小 | 月增量(MB) | 年增量(GB) | TTL |
| ------------------------ | ------------ | -------- | ------------ | ------------ | ---- |
| `student_dashboard_view` | 500K | 200B | 100 | 1.2 | 3 年 |
| `student_errors` | 100K | 150B | 15 | 0.18 | 2 年 |
| `mastery_snapshot` | 250K | 100B | 25 | 0.3 | 5 年 |
| `attendance_logs` | 500K | 80B | 40 | 0.48 | 3 年 |
| `ai_usage_log` | 50K | 120B | 6 | 0.07 | 1 年 |
**总计**:年增量 ~2.3 GB,单节点 ClickHouse 可承载。P6 引入集群时按 class_id hash 分片。
### 11.2 冷热分层
- 热数据:近 6 个月,SSD 存储
- 冷数据:6 个月 - 3 年,HDD 存储(ClickHouse `TTL ... TO VOLUME 'cold'`)
- 归档数据:> 3 年,导出 Parquet 到对象存储(P6 评估)
### 11.3 Kafka topic partition 规划
| Topic | P4 partition | P6 partition | 依据 |
| ---------------------------------------- | ------------ | ------------ | -------------------------- |
| `edu-cdc.next_edu_cloud.core_edu_grades` | 1 | 4 | 高频写入,按 class_id 分区 |
| `edu-cdc.next_edu_cloud.core_edu_exams` | 1 | 2 | 低频 |
| `edu.insight.mastery.updated` | 1 | 4 | 下游多消费者 |
## 12. 数据治理与隐私(v2 新增,长期合规铺垫)
### 12.1 GDPR / 个人信息保护法对齐
| 场景 | 实现 | 阶段 |
| ---------- | ---------------------------------------------------------------------------------------------- | ---- |
| 学生删除权 | ClickHouse 无 UPDATE/DELETE,用 `ALTER TABLE ... DELETE WHERE student_id = ?`(异步 mutation) | P6 |
| 数据导出权 | `GET /analytics/student/{id}/export` → JSON 全量导出(v2 预留 API) | P6 |
| 字段脱敏 | student_id 在日志中 hash 化(structlog processor) | P4 |
| 审计日志 | 所有查询记录 `user_id + query + timestamp` 到 ClickHouse `audit_log` 表(v2 新增表,P6 实现) | P6 |
| 数据保留期 | 见 §11.1 TTL 列 | P4+ |
### 12.2 敏感数据标记
- `student_id`:PII,日志中 hash 化
- `score`:非敏感,可明文
- `attendance_status`:非敏感
- `ai_usage`:含 user_id + prompt_tokens,prompt 内容不入库(ai 服务侧处理)
## 13. 测试策略(v2 新增)
### 13.1 测试金字塔
| 层级 | 工具 | 覆盖范围 | 目标覆盖率 |
| -------- | ----------------------- | -------------------------------------------- | ---------- |
| 单元测试 | pytest + pytest-asyncio | 算法(掌握度计算/预警阈值)、Pydantic 模型 | ≥ 80% |
| 集成测试 | pytest + Testcontainers | ClickHouse DDL + 查询、Kafka producer、Redis | ≥ 60% |
| 契约测试 | pact-python(v2 评估) | gRPC 契约(analytics.proto) | 关键 RPC |
| E2E 测试 | pytest + docker-compose | CDC 链路端到端(MySQL → Debezium → CH) | 关键场景 |
### 13.2 Testcontainers 配置
```python
@pytest.fixture(scope="session")
async def clickhouse_container():
container = ClickHouseContainer("clickhouse/clickhouse-server:24-alpine")
container.start()
yield container
container.stop()
@pytest.fixture(scope="session")
async def kafka_container():
container = KafkaContainer("confluentinc/cp-kafka:7.6.0")
container.start()
yield container
container.stop()
```
### 13.3 测试数据
- 使用 Faker 生成 100 学生 × 10 班级 × 5 科目 × 20 次成绩的合成数据
- 边界用例:空数据、单条数据、超大数据(10000 行)、重复事件(幂等测试)
## 14. Dockerfile 多阶段构建(v2 新增)
```dockerfile
# services/data-ana/Dockerfile
FROM python:3.12-slim AS builder
WORKDIR /app
RUN pip install --no-cache-dir uv
COPY pyproject.toml uv.lock ./
RUN uv sync --frozen --no-dev
FROM python:3.12-slim AS runtime
WORKDIR /app
COPY --from=builder /app/.venv /app/.venv
COPY src/ ./src/
ENV PATH="/app/.venv/bin:$PATH"
ENV PYTHONUNBUFFERED=1
EXPOSE 3006 50055
HEALTHCHECK --interval=30s --timeout=3s --retries=3 \
CMD curl -f http://localhost:3006/healthz || exit 1
CMD ["python", "-m", "uvicorn", "src.data_ana.main:app", \
"--host", "0.0.0.0", "--port", "3006"]
```
**优化点**:
- 多阶段构建,runtime 镜像 < 200MB
- uv 替代 pip,安装速度 10x
- uv.lock 锁定依赖版本
- 健康检查集成
## 15. 完整配置项清单(v2 新增,补全 config.py)
```python
# config.py 完整配置项
class Settings(BaseSettings):
# 服务
service_name: str = "data-ana"
http_port: int = 3006
grpc_port: int = 50055
log_level: str = "INFO"
# ClickHouse
clickhouse_host: str = "edu-clickhouse"
clickhouse_port: int = 8123
clickhouse_user: str = "default"
clickhouse_password: str = ""
clickhouse_database: str = "edu_analytics"
clickhouse_connect_timeout_ms: int = 3000
clickhouse_query_timeout_s: int = 3 # P4 退出标准 5s,查询 3s 超时降级
# Kafka
kafka_brokers: str = "kafka:29092"
kafka_consumer_group: str = "data-ana-cdc"
kafka_cdc_topics: str = (
"edu-cdc.next_edu_cloud.core_edu_grades,"
"edu-cdc.next_edu_cloud.core_edu_exams,"
"edu-cdc.next_edu_cloud.core_edu_homework_submissions,"
"edu-cdc.next_edu_cloud.core_edu_attendance,"
"edu-cdc.next_edu_cloud.classes,"
"edu-cdc.next_edu_cloud.iam_users,"
"edu-cdc.next_edu_cloud.content_knowledge_points"
)
kafka_mastery_topic: str = "edu.insight.mastery.updated"
kafka_warning_topic: str = "edu.insight.warning.triggered"
kafka_ai_usage_topic: str = "edu.insight.ai.usage"
kafka_enable_auto_commit: bool = False # v2: 手动 commit
kafka_auto_offset_reset: str = "latest"
# iam gRPC
iam_grpc_endpoint: str = "iam:50052"
iam_grpc_timeout_s: int = 2
datascope_cache_ttl_s: int = 300
# Redis
redis_url: str = "redis://edu-redis:6379/3"
redis_pool_size: int = 10
redis_socket_timeout_ms: int = 200
# OTel
otel_endpoint: str = "http://otel-collector:4317"
otel_service_name: str = "data-ana"
# 掌握度算法
mastery_window_size: int = 5
mastery_decay_base: float = 0.6
mastery_forgetting_half_life_days: int = 30
mastery_min_samples: int = 3
# 预警阈值
warning_low_mastery_threshold: float = 0.4
warning_critical_mastery_threshold: float = 0.2
warning_score_drop_percent: float = 0.2
warning_absent_per_week: int = 3
# 降级
degraded_mode_enabled: bool = True
model_config = SettingsConfigDict(env_file=".env", env_prefix="DATA_ANA_")
```
## 16. API 版本化策略(v2 新增)
### 16.1 HTTP 版本化
- URL 前缀:`/analytics/v1/...`(当前省略 v1,v2 起强制)
- 破坏性变更:新增 `/analytics/v2/...`,旧 v1 保留 6 个月
- 非破坏性变更:直接在 v1 上加字段(Pydantic 模型字段可选)
### 16.2 gRPC 版本化
- proto package:`next_edu_cloud.analytics.v1`(当前)
- 破坏性变更:新增 `next_edu_cloud.analytics.v2`,v1 保留
- buf breaking 检查:CI 强制(FILE 级别)
### 16.3 事件版本化
- 事件 schema 字段:`schema_version: "v1"`(必填)
- 破坏性变更:发布 `v2` 事件,consumer 同时支持 v1 + v2(6 个月过渡)
- 非破坏性变更:直接加字段,schema_version 不变
## 17. Projection 演进路径(v2 新增)
### 17.1 当前(P4):CDC 直连
```
MySQL binlog → Debezium → Kafka (edu-cdc.*) → data-ana CdcConsumer → ClickHouse
```
- 优点:实时性好(< 5s),无需 core-edu 改造
- 缺点:依赖 binlog schema,core-edu 表结构变更可能破坏 CDC
### 17.2 未来(P4 后期评估):双消费
```
MySQL binlog → Debezium → Kafka (edu-cdc.*) → data-ana CdcConsumer → ClickHouse(宽表)
↓
core-edu Outbox → Kafka (edu.domain.events) → data-ana EventConsumer → ClickHouse(业务事件视图)
```
- CDC 通道:保持宽表实时性
- 领域事件通道:业务语义清晰,schema 稳定(proto 定义)
- 两通道并存,分别写不同 ClickHouse 表,避免重复
### 17.3 远期(P6):CDC + CQRS 物化视图
```
CDC → ClickHouse 原始表 → MaterializedView → 宽表(自动聚合)
```
- 减少应用层计算,由 ClickHouse MV 自动维护宽表
- 掌握度计算仍需应用层(复杂算法 CH 难以表达)
## 18. 未来阶段铺垫(v2 新增,长远架构)
### 18.1 P5 铺垫(AI 网关 + 沟通通知)
| 演进点 | P4 准备 | P5 实现 |
| ---------------- | ----------------------------------------------------------- | --------------------------------------------- |
| AI 用量消费 | `ai_usage_log` 表已建 + `edu.insight.ai.usage` topic 已订阅 | ai 服务发布事件,data-ana 落库 + 用量统计 API |
| gRPC stream 推送 | `SubscribeMasteryUpdate` RPC 已在 proto 提案(§4.2) | 实现 server-streaming,AI 实时获取掌握度变化 |
| 预警触发通知 | `WarningTriggered` 事件已设计(§5.2) | msg 服务消费事件 → 推送 push-gateway |
### 18.2 P6 铺垫(生产硬化)
| 演进点 | P4 准备 | P6 实现 |
| ----------------- | --------------------------------------------- | ------------------------------------ |
| CDC 水平扩展 | ExamCache 抽象层 + Redis 配置已就绪(§5.3) | 多实例 + HPA + partition 扩容 |
| ClickHouse 集群 | DDL 用 ReplacingMergeTree(兼容 Replicated*) | ReplicatedMergeTree + 副本 + 分片 |
| Service Mesh mTLS | gRPC server 已抽象 interceptor | Istio 注入 sidecar,mTLS 自动双向 |
| 数据治理 | student_id 日志 hash + TTL 设计(§12) | ALTER TABLE DELETE + 审计日志表 |
| 配置中心 | Settings 用 pydantic-settings(支持热更新) | Consul 接入,watch 配置变更 |
| 灾备 | Kafka producer idempotent + transactional | Kafka 多副本 + ClickHouse 跨机房副本 |
### 18.3 远期(P7+)预留
- **实时流式分析**:Flink 替代 CdcConsumer,支持窗口聚合(如"近 7 天平均掌握度"实时计算)
- **数据湖**:ClickHouse 冷数据 → S3/MinIO + Iceberg 表格式,支持 ad-hoc 查询
- **AI 驱动的学情诊断**:data-ana 提供数据,ai 服务提供模型,组合成"智能诊断报告"
- **多租户**:当前 school_id 隐含在 class_id 中,远期可显式多租户隔离
## 19. 服务审计表 v2(黄金模板对齐)
| 维度 | 状态 | 说明 |
| ------------- | ---- | ----------------------------------------------------------------------- |
| 契约(proto) | ✅ | analytics.proto 已定义 3 RPC,v2 提案扩展 7 RPC |
| 路由 | ✅ | HTTP 13 端点 + gRPC 10 RPC |
| 数据访问 | ✅ | ClickHouse 5 宽表 + Redis 缓存 + iam gRPC |
| 鉴权 | ✅ | FastAPI Depends 等价物 + 7 权限点 |
| 错误处理 | ✅ | 13 错误码 `DATA_ANA_*` + ActionState 信封 |
| 可观测 - 日志 | ✅ | structlog + JSONRenderer(v2 改进) |
| 可观测 - 指标 | ✅ | 16 Prometheus 指标 |
| 可观测 - 链路 | ✅ | OTel + gRPC interceptor + CDC span 注入 |
| 健康检查 | ✅ | /healthz + /readyz(v2 补 redis/iam_grpc) |
| 配置管理 | ✅ | pydantic-settings 完整配置项(§15) |
| 降级模式 | ✅ | 4 降级场景(CH/Kafka/Redis/iam) |
| 响应信封 | ⚠️ | v2 设计对齐 ActionState,实现阶段需重构(P0 整改) |
| Redis | ⚠️ | v2 设计就绪,实现阶段需新增 redis-py 依赖 + RedisClient |
| gRPC server | ⚠️ | v2 设计就绪,实现阶段需新增 grpc.aio + betterproto + server interceptor |
| Dockerfile | ⚠️ | v2 设计多阶段,实现阶段需新增 |
| 测试 | ⚠️ | v2 策略就绪,实现阶段需补 pytest + Testcontainers |
| CDC 幂等性 | ⚠️ | v2 改进手动 commit + Redis SETNX,实现阶段需重构 |
| CDC 水平扩展 | ✅ | v2 设计 P6 演进路径,P4 单实例可工作 |
| 数据治理 | ✅ | v2 设计 GDPR 对齐方案,P6 实现 |
| 容量规划 | ✅ | v2 估算年增量 2.3GB,单节点可承载 |
---
**AI Agent**: ai11 (data-ana)
**Branch**: 单仓库并行模式(直接提交 main)
**Coordinator**: coord-ai
**v2 修订依据**: ai-allocation.md §3.2 重新分配 + 004 §1.2/§4.2/§7.2/§11.5/§12.2/§15.3 + ai-allocation §5 设计重点 + pending-features P4 退出标准
**v2 审核结论**: 16 项遗漏已全部补强(3 项 P0 + 5 项 P1 + 8 项 P2),文档进入实现阶段无阻塞