Router 数据回流管道(数据飞轮)设计提案
状态:阶段 1/2 已实施(2026-09)。当前线上形态:wrapper 直报 → router 摄入 → PG(计费/明细同步路径不动)+ 脱敏后分流 Kafka(栈内自建 KRaft)→ ① worker 自消费攒批写 OSS 语料冷层(
lake/<topic>/dt=YYYY-MM-DD/hour=HH/part-*.jsonl.gz)② 第二路消费写 ClickHouse(栈内自建单节点)→ console 实时监控页 + 飞书告警。阶段 3 为设计目标。基础架构见 architecture-router.md。
1. 目标与边界
数据飞轮服务三类消费,优先级与 SLA 各不相同——这是管道设计的第一性原理:
| 消费方 | 数据内容 | 时效要求 | 可靠性要求 | 现状 |
|---|---|---|---|---|
| 计费/限额(rated events) | 指标行(tokens/成本/归属) | 秒级(欠费拦截) | 最高:不丢不重(幂等键) | ✅ PG 同步落库 |
| 决策明细存证 | 指标行 + 错误分类 | 分钟级可查 | 高:永久保留 | ✅ PG 月分区 |
| 实时看板/告警 | 指标行聚流 | 分钟级(1-5min) | 中:允许重放补齐 | ❌ 只有天级 rollup |
| 训练语料湖 | prompt 文本(脱敏后) | 小时/天级 | 中:批量不丢 | ⚠️ PG 采样表 30 天 TTL |
| 标注工作台 | 语料子集 + 标注 | 交互式 | 高:随机读写 | ✅ PG |
核心取舍:计费/明细链路不走异步管道——账务准确性要求同步、幂等、强一致,现有「摄入即计价落 PG」路径保留不动。异步管道承载观测(看板/告警)与语料——允许 at-least-once、秒级延迟、重放补齐。两条链路共用同一个事件源,但在 router 摄入侧分流。
2. 目标架构
┌──────────────────────────────────────────────┐
│ LiteLLM wrapper(veyra_hooks,不动) │
│ fire-and-forget 批量上报 StandardLoggingPayload│
└──────────────────┬───────────────────────────┘
▼
┌──────────────────────────────────────────────┐
│ router api(2 副本,无状态摄入,现有) │
│ ① PG 同步落库:rated 计价 + 明细 + 限额计数 │← 计费链路(同步、幂等)
│ ② 脱敏(PII/密钥模式) │
│ ③ 转发 Kafka(摄入批 → produce,本地磁盘 │
│ 缓冲队列兜底,Kafka 故障不阻塞摄入) │← 观测/语料链路(异步)
└──────────────────┬───────────────────────────┘
▼
┌───────────────────────┐
│ Kafka(阶段 1 栈内自建 │
│ KRaft 单节点;阶段 2 │
│ 起评估迁托管) │
│ topic: inference-events│
└───────────┬───────────┘
┌─────────────────────┼──────────────────────────┐
▼ ▼ ▼
┌──────────────────┐ ┌────────────────────┐ ┌────────────────────┐
│ Flink 实时计算 │ │ Kafka Connect/OSS │ │ 告警引擎 │
│ - 二次清洗/增强 │ │ 投递:原样批量落 OSS │ │ (SLS 告警或 Flink │
│ - 窗口聚合 1m/5m │ │ JSONL/Parquet │ │ CEP 规则) │
│ - 分类器预标 │ │ dt=/org=/model= 分区 │ │ → 钉钉/webhook │
└────────┬─────────┘ └─────────┬──────────┘ └────────────────────┘
▼ ▼
┌──────────────────┐ ┌────────────────────┐
│ ClickHouse(托管) │ │ OSS 语料湖(冷层) │
│ 分钟级聚合物化视图 │ │ 训练原料长期保存 │
│ → console 实时看板 │ │ 训练快照不可变版本 │
└──────────────────┘ └────────────────────┘
▲
┌─────────────────────────────────┘
│ 标注工作台(PG,现有 prompt_samples)
│ 标注完成 → 导出不可变训练快照到 OSS(sample_id 关联原料层)2.1 分流点:router 摄入侧转发 Kafka
wrapper 不感知 Kafka(不新增依赖、不改协议):router 摄入批在 PG 落库成功后,异步 produce 到 Kafka。
- 脱敏在转发前完成(正则 + 规则:密钥串/手机号/邮箱/身份证模式 → 打码)——语料湖里的文本永远不含明文敏感信息;原始明文只活在 PG 采样热层(30 天 TTL)
- Kafka 故障兜底:produce 失败写本地磁盘 WAL 文件(或内存环形队列 + 磁盘溢出),后台重放;绝不阻塞摄入主路径
- 幂等键复用:事件带
event_id,下游 Flink 按 event_id 去重(与我方摄入幂等键同源)
2.2 实时看板层:Flink + ClickHouse
- Flink 作业:1min/5min 滚动窗口聚合(org × model × status × error_class × key 维度),输出 ClickHouse
- ClickHouse 物化视图链路:明细表(Kafka 引擎表)→ 聚合表(SummingMergeTree)→ console 看板查询
- 看板内容:实时 QPS/错误率/TTFT P50/费用速率,按 org/key/模型/供应商切片;console 新增「实时监控」页(或直接 Grafana 嵌入链接——iframe 禁令仅针对终端门户,运营台可用)
- 告警:两类规则——阈值类(5min 上游错误率 >10%、某供应商不可用、费用速率突增)走 SLS 告警或 Flink CEP → 钉钉机器人/webhook;漂移类(成本漂移、价格同步待办)保持现有 DB 检测 + console 露出
2.3 语料湖与标注闭环
- 原料层:Kafka Connect(或 Flink Sink)批量落 OSS,Parquet/JSONL,
dt=yyyy-MM-dd/hour=HH/org=.../model=.../分区——批量文件,不逐条写对象(小对象塌方规避) - 标注层:PG
prompt_samples不动(随机读写强需求);标注结果以sample_id关联原料层 - 训练快照:标注完成/周期触发 → join 原料 + 标注 → 不可变快照(
snapshots/v<YYYYMMDD-HHmm>/)落 OSS——训练可复现的版本化真相源 - 早期不上 Iceberg/Hudi(表格式治理是下一阶段的事),前缀约定 + manifest.json 足够
2.4 阿里云组件选型
| 角色 | 选型 | 备选 | 理由 |
|---|---|---|---|
| 消息缓冲 | 栈内自建 KRaft Kafka 单节点(deploy/router-prod/kafka.yaml);阶段 2 起评估迁托管 Kafka | 托管 Kafka Serverless / 日志服务 SLS | 托管最小预留 300MB/s(~8700 元/年)对 KB/s 级实际流量是数百倍过剩;自建单点由摄入 WAL + 消费 offset 纪律 + PG 真相源回填兜底(详见 deploy/router-prod/README「数据回流」);托管仍是阶段 2 的归宿(改 bootstrap + 双跑追平即可迁) |
| 流计算 | 不做(CH 实时 SQL 覆盖窗口聚合) | 实时计算 Flink 版 | 当前量级直查明细足够;Flink 推迟到有 CEP 复杂序列告警需求时 |
| 实时 OLAP | 栈内自建 ClickHouse 单节点(deploy/router-prod/clickhouse.yaml);规模上来再迁云数据库 ClickHouse | 托管 ClickHouse 社区版 / AnalyticDB | 托管最低规格月费数百元对 2 条/s 过剩(同 Kafka 决策逻辑);CH 是观测副本,挂了不影响摄入/账务 |
| 对象存储 | OSS(标准-低频分层) | — | 生产本在阿里云,内网 endpoint 免流量费 |
| 告警 | alert_events(PG)+ 飞书群机器人加签推送(自研 60s 评估循环,阈值 runtime_config 热调) | SLS 告警 / Flink CEP | 三条规则(通道不可用/错误率突增/费用速率)定时查询足够;CEP 留待复杂序列规则 |
| 网络 | 全部走 VPC 内网 endpoint;k3s 节点与中间件同 VPC | — | 免公网流量费 + 安全 |
规模估算:500 万调用/月 ≈ 均值 2 条/s、峰值 ~50 条/s、单事件 ~2KB → 自建单节点 Kafka(250m/512Mi 起步)绰绰有余;阶段 2 的 Flink 1-2CU、ClickHouse 单副本才是按需花钱的项。
3. 分阶段落地路径
每一步独立可验收、可回滚,不为终态提前买单:
阶段 0(现状,已完成):PG 明细永久分区 + rated 计费 + 天级 rollup + 采样语料 PG 热层。
阶段 1 — 事件流化 + 语料冷层(收益:语料长期化 + 多消费者底座)——已实施
- router 摄入侧 Kafka producer(
app/services/event_stream.py:落库成功后转发;发送失败落 pod 本地 WALkafka-wal.jsonl,60s 后台重放)+ 脱敏器(app/services/redact.py,转发前强制过 scrub) - 消息层为栈内自建 KRaft Kafka 单节点(
deploy/router-prod/kafka.yaml,topic auto-create 3 分区);托管 Kafka 起步价决策记录见 §2.4 选型表 - 投递方式调整为 router-worker 自消费攒批写 OSS(
app/services/lake_writer.py:满 500 条或 60s 冲刷;offset 上传成功才 commit,at-least-once)——替代原 Kafka Connect 托管流出方案,零额外组件 - 对象布局:
lake/<topic>/dt=YYYY-MM-DD/hour=HH/part-<millis>-<seq>-<uuid8>.jsonl.gz;topic:inference-events(决策指标行,无明文)/inference-samples(脱敏后语料) - 验收:OSS 上有按天分区文件,与 PG 明细对账条数一致(±去重窗口)
阶段 2 — 实时看板 + 告警(收益:分钟级运营可见性)——已实施(调整版)
- 栈内自建 ClickHouse 单节点(
deploy/router-prod/clickhouse.yaml)存明细router.events(MergeTree,TTL 90 天);写入为 router-worker 第二路消费者ch_writer(group=router-ch-writer,攒批 500 条/60s HTTP JSONEachRow 写入, offset 写入成功才 commit)——替代原 Kafka 引擎表/Connect 方案,与 lake_writer 对称 - 无 Flink:CH 实时 SQL(
toStartOfMinutegroup by +uniqExact(event_id)去重) 覆盖窗口聚合——当前量级直查明细即毫秒级;Flink 推迟到有 CEP 复杂序列告警需求时 - console 新增「实时监控」页(echarts 分钟级曲线 + 近 5min 通道状态 + 告警时间线; 不引 Grafana,单入口)
- 告警走 alert_events(PG)+ 飞书群机器人加签推送(替代 SLS/CEP): 通道不可用 critical / 上游错误率突增 warning / 费用速率异常 warning, 60s 评估 + 30min 冷却 + 阈值 runtime_config 热调
- 验收:端到端延迟 < 2min(调用发生 → 看板可见);告警演练命中
阶段 3 — 训练闭环(收益:飞轮完整)
- 标注快照导出器(PG 标注 join OSS 原料 → 不可变快照)
- 分类器重训管线消费快照;产物注册表灰度(现有状态机)
- 验收:某版本产物由快照 X 训练 → manifest 可溯源到快照路径
4. 明确不做 / 风险
- 计费永不走 Kafka:异步管道允许 at-least-once,账务不允许;rated events 保持摄入即落 PG
- wrapper 不直写 Kafka:保持单点出口(router),wrapper 协议不变
- 不上 Iceberg/Hudi(早期):前缀约定 + manifest 足够;等需要 ACID 演变/快照隔离再上表格式
- 脱敏是阶段 1 的前置门禁:不脱敏不进 Kafka/OSS——这是合规红线,不是可选项
- 双写窗口:阶段 1 后 PG 明细与 OSS 原料短暂双写——以 PG 为准对账;阶段 3 后明细可评估瘦身(PG 热查询窗口 + OSS 冷存证),但那是后话,本文不承诺