Skip to content

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 去重(与我方摄入幂等键同源)
  • 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 本地 WAL kafka-wal.jsonl,60s 后台重放)+ 脱敏器(app/services/redact.py,转发前强制过 scrub)
  • 消息层为栈内自建 KRaft Kafka 单节点deploy/router-prod/kafka.yaml,topic auto-create 3 分区);托管 Kafka 起步价决策记录见 §2.4 选型表
  • 投递方式调整为 router-worker 自消费攒批写 OSSapp/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(toStartOfMinute group 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 冷存证),但那是后话,本文不承诺

基于内网部署的企业级 AI 智能体平台