Day 40 Hard ML Infra Feature Store Serving Skew A/B

特征平台与 ML 基础设施 — 训练与推理共用一份特征的工程学Feature Platform & ML Infra: Dual Store, Point-in-Time Correctness, Model Serving, Experimentation

问题场景 + 需求约束

设计一个支撑外卖/风控级别的特征平台:100+ 模型(配送 ETA、门店排序、欺诈识别)共用一套特征,在线推理峰值 1000 万次特征读/秒、单次预测要取 200 个特征、p99 < 10ms;同时离线要能为几十亿样本回溯出「当时那一刻」的特征做训练。这就是 DoorDash、Uber 真实的量级。

核心矛盾:训练用离线批处理(吞吐优先),推理用在线低延迟(延迟优先),但两边必须是同一个特征值——否则模型线下 AUC 0.85、线上掉到 0.79,这就是臭名昭著的 training-serving skew。平台要解决三件事:存哪、怎么保证一致、怎么服务与验证

高层架构

graph LR
    SRC["数据源
事件流 / 数仓 / DB CDC"] subgraph 计算层 BATCH["批计算
Spark / SQL"] STREAM["流计算
Flink"] end REG["特征注册表
定义 + 版本 + 血缘"] OFF[("离线 Store
Parquet / 数仓")] ON[("在线 Store
Redis / DynamoDB")] TRAIN["训练
PIT join 取样本"] SERVE["模型服务
取特征 + 推理"] SRC --> BATCH & STREAM REG -.约束定义.-> BATCH & STREAM BATCH --> OFF BATCH -->|物化 materialize| ON STREAM --> ON OFF --> TRAIN ON --> SERVE classDef src fill:#1a2530,stroke:#64c8ff,color:#e8eef5 classDef comp fill:#0e2030,stroke:#5eead4,color:#e8eef5 classDef store fill:#2a1530,stroke:#ff7ab6,color:#e8eef5 classDef sink fill:#1a1a30,stroke:#ffb450,color:#e8eef5 class SRC src class BATCH,STREAM,REG comp class OFF,ON store class TRAIN,SERVE sink

同一份特征定义(注册表)驱动批与流两条管道,分别落离线/在线两个 Store,喂给训练与推理

组件职责注册表是唯一真相源,声明特征的实体键、类型、转换逻辑与新鲜度 SLA;批计算做重回溯与 T+1 特征、流计算做秒级新鲜特征;离线 Store按时间分区保留全历史(供 point-in-time 取数),在线 Store只存每个实体的最新值(供毫秒读)。这就是业界共识的双 Store(dual-store)架构。

关键技术点

1. 双 Store:离线仓库 vs 在线 KV,一个定义两处物化

原理:训练要「几十亿样本 × 每样本回溯历史特征」,是吞吐型全表扫描,天然属列存数仓(Parquet/Iceberg,按 event_ts 分区)。推理要「给定一个用户,10ms 内取 200 个特征」,是点查型随机读,天然属内存 KV。同一个特征在两个引擎各物化一份,靠注册表保证语义一致。

Trade-off(在线 Store 选型):
# 在线取数:一次 pipeline 批量拉 200 个特征键,避免 200 次 RTT
def fetch_online(entity_ids, feature_names):
    keys = [f"{fn}:{eid}" for eid in entity_ids for fn in feature_names]
    vals = redis.mget(keys)          # 单次往返,200 键 ~1ms
    return assemble(entity_ids, feature_names, vals)

# 物化 materialization:批管道把最新值刷进在线 Store
# 只写「每实体最新」,不保留历史 → 在线 Store 体积可控
现实案例:

2. Point-in-Time 正确性:杜绝标签泄漏的核心机制

原理:训练样本是「(实体, 标签时刻 T, 标签)」,为它拼特征时,只能用 event_ts ≤ T 的最新特征值。若图省事直接 join 特征表的当前值,就把「未来」信息喂给了模型——线下指标虚高、线上崩盘。正确做法是 point-in-time join(as-of join):对每个样本按实体键找「T 之前最近一条」特征。

Trade-off(保证一致的两条路线):
# Point-in-time join 伪代码(as-of join,防标签泄漏)
# labels: (entity_id, label_ts, y)   features: (entity_id, event_ts, value)
for row in labels:
    cand = features[entity_id == row.entity_id
                    and event_ts <= row.label_ts]   # 严格 ≤ 标签时刻
    row.feature = cand.sort_by(event_ts).last()      # 取当时最新的一条
# 引擎实现:Spark 按 entity_id 分区 + event_ts 排序做 merge as-of join
现实案例:

3. 模型服务:特征拉取 + 推理的延迟预算切分

原理:一次线上预测的 p99 预算(如 30ms)要在「取特征 + 前处理 + 模型前向 + 后处理」间分配。经验上取特征常占一半以上延迟——200 个特征散在多个实体、多个 Store。优化点:批量 pipeline 取数、把同请求的特征并发拉取、对高频实体做本地 L1 缓存。

Trade-off(服务形态):
现实案例:

4. A/B 实验平台:让「模型上线」变成可度量的因果实验

原理:新模型不能因线下 AUC 涨就上线——线下指标和业务指标(下单率、GMV)常背离。要用随机分流的在线实验建立因果:一致性哈希把用户稳定分到 control/treatment,跑够样本量后比较业务指标的统计显著性,同时盯 guardrail 指标(延迟、错误率、退款率)防止「主指标涨、副作用炸」。

Trade-off(关键设计点):
# 稳定分桶:同一 user 永远落同一组(可复现、无闪烁)
def bucket(user_id, exp_key, treat_ratio=0.5):
    h = hash_u64(f"{exp_key}:{user_id}")   # 加实验名做盐,各实验独立正交
    return "treatment" if (h % 10000) / 10000 < treat_ratio else "control"
现实案例: Netflix、Uber、Airbnb 均自建实验平台,核心都是「稳定分桶 + guardrail 指标 + 方差削减」。Airbnb 的 Chronon 特征定义可直接被实验消费,做到「新特征→上模型→A/B」闭环。

扩展与优化

常见陷阱 + 面试问题

1. 训练用当前值 join 特征表:最经典的标签泄漏。必须 point-in-time join,用 event_ts ≤ label_ts 严格约束。
2. 训练与服务两套转换代码:Python 版和 Java 版对同一特征算出不同值 = skew。共享一份定义/DSL 是根治之道。
3. 把在线 Store 当真相源:Redis 宕机 RPO 有损,它是最新值的缓存,历史与真相在离线 Store 与数仓。
4. A/B 无 guardrail:只看主指标涨就全量,结果延迟翻倍、退款率飙升。必须同时监控副作用指标。

深入资源

深入思考(点击展开答案)

1. 一个特征在离线算出来是 3.7,在线算出来是 3.5,模型没报错但线上效果掉了 5%。你怎么系统性地定位并根治这类「无声 skew」?

定位分三层

  • 数据层:对同一批 (实体, 时刻) 分别从离线 PIT 取数和在线打日志取数,逐特征对比分布与逐条差值,找出偏差最大的特征。
  • 逻辑层:偏差大的特征,检查批与流是否用了两套实现(如时间窗口口径不同:离线 [T-7d, T) 左闭右开,在线不小心含了当前事件)、类型/精度(float32 vs float64)、缺失值填充策略不一致。
  • 时间层:在线特征可能有物化延迟——训练用了 T 时刻的值,但线上此刻 Store 里还是 T-5min 的旧值,等价于隐性 skew。

根治:① 特征转换用同一份代码/DSL(Chronon 思路),杜绝双实现;② 生产落地持续 skew 监控——每天抽样对比在线日志 vs 离线回溯,超阈值告警;③ 新特征优先用 log-and-wait,让训练数据天生等于服务数据。Google Rules of ML 的核心建议正是「log serving features,train on them」。

2. 为什么不能用一个统一的存储引擎同时满足训练的全表扫描和推理的毫秒点查?从存储物理层解释这个「被迫双 Store」。

本质是访问模式对存储布局的要求相互冲突

  • 训练:几十亿行 × 少数列的顺序扫描 + as-of join。最优布局是列存 + 按时间分区(Parquet/Iceberg),压缩比高、只读需要的列、按分区裁剪。但它点查一行极慢(要扫 row group)。
  • 推理:给一个 key 取一行的多个字段,要 O(1) 随机读。最优是内存哈希表/LSM 点查(Redis/KV),但它全表扫描要遍历全部 key,且只存最新值、无历史。

列存扫描友好但点查慢,KV 点查友好但扫描/历史差——没有单一引擎能同时把两条曲线都压到最优。所以工程上把同一份特征物化两份:历史全量进列存供训练,最新快照进 KV 供推理,代价是要维护两者的一致性(正是双 Store 架构存在的理由)。这和 Day 20 讲的「同一份数据派生多种存储视图」是同一个思想。

3. 双边市场(外卖:用户/骑手/商家)里,把「新的骑手调度模型」按用户随机分 A/B 会得到偏乐观还是偏悲观的结论?为什么?该换什么实验设计?

通常偏乐观(高估收益)。因为 treatment 与 control 用户共享同一个骑手池:新模型如果把骑手更高效地调给 treatment 用户,等于从 control 用户那里「抢」了运力,人为拉大了两组差距——你测到的不是「新模型的绝对收益」,而是「treatment 挤占 control 后的相对差」。全量上线后大家都在同一池子里,收益会缩水甚至消失(SUTVA 假设被违反)。

换设计:用切换实验(switchback / 区域-时间随机)——按「城市 × 时间片」为单位整体切 control/treatment,让同一时空里所有人用同一策略,供需关系闭合在实验单元内,消除跨组干扰。代价是实验单元变少、方差变大、需要更长周期,且要处理时间自相关。Uber/DoorDash 的调度类实验普遍用 switchback。

4. 在线 Store 里某顶流门店的特征被 QPS 100 万单点打爆(联想 Day 2 的热 key)。特征平台场景下,哪些缓解手段可用、哪些不可用,为什么?

先套用 Day 2 的热 key 工具箱,再按特征场景筛:

  • ✅ app 进程内 L1 缓存:最有效。特征在一个推理周期内基本不变,本地缓存 TTL 几秒即可把 Redis 压力降几个数量级——DoorDash 客户端缓存正是此法,提升约 70%。
  • ✅ 多副本读:热特征复制到多个 Redis 节点分摊读。特征读多写少,副本 lag 影响小,很适用。
  • ⚠️ Key 分片(拆成 N 个子 key 求和):只对可累加的计数型特征成立(如「门店今日订单数」);对 embedding、比率类特征无法拆合,不可用。
  • ✅ 请求内去重/合并:同一请求批里多个候选都引用同一门店特征时,只取一次。
  • ❌ 直接调大 TTL 换新鲜度:风控/调度特征对新鲜度敏感,牺牲新鲜度换命中率可能直接损害模型效果,要按特征的新鲜度 SLA 分别定策略。

实战组合:本地 L1(吸收绝大部分读)+ 热特征多副本(兜底),且缓存 TTL 严格对齐该特征的新鲜度要求。

5. 团队想「立刻」用一个全新特征训练模型,但它在线上还没积累任何日志。log-and-wait 与 backfill 两条路你怎么选?各自的隐藏风险是什么?

想立刻训练 → 只能 backfill(log-and-wait 要等日志积累,快则数天慢则数周)。但 backfill 的隐藏风险是「历史无法完美重建」

  • 数据不可回溯:若该特征依赖某个只存最新值、不留历史的上游(如某维度表被原地覆盖),你根本无法算出「三个月前那一刻」的值,回溯出来的是被污染的近似。
  • 回溯逻辑 ≠ 在线逻辑:backfill 用批代码重算,一旦和未来的在线流代码有口径差异,训练出的模型上线即 skew。
  • 幸存者/时点偏差:回溯常用「现在还存在的实体」,把已流失用户漏掉,样本分布偏移。

务实策略:先 backfill 快速验证特征是否有信号(离线跑通、看重要性),同时立刻开启在线打日志;一旦日志够了,切换到 log-and-wait 训练做「正式版」,用它对齐线上真实分布。即backfill 探路、log-and-wait 定稿。切忌用 backfill 版直接长期迭代——skew 会悄悄累积。