设计一个支撑外卖/风控级别的特征平台: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)架构。
原理:训练要「几十亿样本 × 每样本回溯历史特征」,是吞吐型全表扫描,天然属列存数仓(Parquet/Iceberg,按 event_ts 分区)。推理要「给定一个用户,10ms 内取 200 个特征」,是点查型随机读,天然属内存 KV。同一个特征在两个引擎各物化一份,靠注册表保证语义一致。
# 在线取数:一次 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 体积可控
原理:训练样本是「(实体, 标签时刻 T, 标签)」,为它拼特征时,只能用 event_ts ≤ T 的最新特征值。若图省事直接 join 特征表的当前值,就把「未来」信息喂给了模型——线下指标虚高、线上崩盘。正确做法是 point-in-time join(as-of join):对每个样本按实体键找「T 之前最近一条」特征。
# 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
原理:一次线上预测的 p99 预算(如 30ms)要在「取特征 + 前处理 + 模型前向 + 后处理」间分配。经验上取特征常占一半以上延迟——200 个特征散在多个实体、多个 Store。优化点:批量 pipeline 取数、把同请求的特征并发拉取、对高频实体做本地 L1 缓存。
原理:新模型不能因线下 AUC 涨就上线——线下指标和业务指标(下单率、GMV)常背离。要用随机分流的在线实验建立因果:一致性哈希把用户稳定分到 control/treatment,跑够样本量后比较业务指标的统计显著性,同时盯 guardrail 指标(延迟、错误率、退款率)防止「主指标涨、副作用炸」。
# 稳定分桶:同一 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"
event_ts ≤ label_ts 严格约束。定位分三层:
[T-7d, T) 左闭右开,在线不小心含了当前事件)、类型/精度(float32 vs float64)、缺失值填充策略不一致。根治:① 特征转换用同一份代码/DSL(Chronon 思路),杜绝双实现;② 生产落地持续 skew 监控——每天抽样对比在线日志 vs 离线回溯,超阈值告警;③ 新特征优先用 log-and-wait,让训练数据天生等于服务数据。Google Rules of ML 的核心建议正是「log serving features,train on them」。
本质是访问模式对存储布局的要求相互冲突:
列存扫描友好但点查慢,KV 点查友好但扫描/历史差——没有单一引擎能同时把两条曲线都压到最优。所以工程上把同一份特征物化两份:历史全量进列存供训练,最新快照进 KV 供推理,代价是要维护两者的一致性(正是双 Store 架构存在的理由)。这和 Day 20 讲的「同一份数据派生多种存储视图」是同一个思想。
通常偏乐观(高估收益)。因为 treatment 与 control 用户共享同一个骑手池:新模型如果把骑手更高效地调给 treatment 用户,等于从 control 用户那里「抢」了运力,人为拉大了两组差距——你测到的不是「新模型的绝对收益」,而是「treatment 挤占 control 后的相对差」。全量上线后大家都在同一池子里,收益会缩水甚至消失(SUTVA 假设被违反)。
换设计:用切换实验(switchback / 区域-时间随机)——按「城市 × 时间片」为单位整体切 control/treatment,让同一时空里所有人用同一策略,供需关系闭合在实验单元内,消除跨组干扰。代价是实验单元变少、方差变大、需要更长周期,且要处理时间自相关。Uber/DoorDash 的调度类实验普遍用 switchback。
先套用 Day 2 的热 key 工具箱,再按特征场景筛:
实战组合:本地 L1(吸收绝大部分读)+ 热特征多副本(兜底),且缓存 TTL 严格对齐该特征的新鲜度要求。
想立刻训练 → 只能 backfill(log-and-wait 要等日志积累,快则数天慢则数周)。但 backfill 的隐藏风险是「历史无法完美重建」:
务实策略:先 backfill 快速验证特征是否有信号(离线跑通、看重要性),同时立刻开启在线打日志;一旦日志够了,切换到 log-and-wait 训练做「正式版」,用它对齐线上真实分布。即backfill 探路、log-and-wait 定稿。切忌用 backfill 版直接长期迭代——skew 会悄悄累积。