特征平台
InfoQ《5 年迭代 5 次》写的是推荐特征生产,不是算法公式。早期 MapReduce/Storm 各搞各的 profile;之后上 Flink;再迁 Flink SQL;再变成「以离线计算为核心」的统一架构。
滑动 / 完播日志
BMQ
Flink 窗口聚合
Abase 特征 KV
精排点查
同一条队列还会分流到湖仓,做离线 / 在线训练
轻在线,重离线
旧方案把明细按时间片塞进 KV,请求来了在线聚合窗口。新方案把窗口聚合放在 Flink State(RocksDB,吃本地 SSD)完成,把结果推到 Abase 一类在线 KV。在线特征服务只做点查。否则晚高峰数百万 QPS 的 Feed 会把存储当成计算引擎打爆。
特征类型
- 无状态 ETL:过滤即可得到的属性。
- 有状态窗口:最近 1h 点赞、session 看播时长。
- 序列:最近 100 次曝光。
- 图特征:二跳关系,例如看的最多的主播收到最多的礼物。
数据源抽象成 Schema Table,支持 Window Join、Interval State Join、Lookup Join(Abase、RPC、Hive)。
规模
日均 PB;特征时效要求分钟级。直播/电商 Flink 作业状态公开到 60TB,并预期单作业 100TB。为此做了 State Cache(减序列化,CPU 约降 50%)和 PB IDL 裁剪(大 Topic 上百字段只用几个,CPU 约降 30%)。
样本与穿越
短视频主标签是完播、停留、有效播放,点击会把封面党训进模型。负例来自曝光未播、快速划走。训练时若用了请求之后才发生的统计,线上没有这个未来值,离线 AUC 会虚高。特征平台要保证训练与 serving 同一份定义,并用请求发生时刻的 snapshot join。
Join 发生在流上:曝光日志与后验动作按 uid+item+request_id 在 Flink 里等一个延迟窗口。窗口外的样本进离线补训。