learntoall

特征平台

InfoQ《5 年迭代 5 次》写的是推荐特征生产,不是算法公式。早期 MapReduce/Storm 各搞各的 profile;之后上 Flink;再迁 Flink SQL;再变成「以离线计算为核心」的统一架构。

滑动 / 完播日志
BMQ
Flink 窗口聚合
Abase 特征 KV
精排点查

同一条队列还会分流到湖仓,做离线 / 在线训练

流程图 · 日志进队列,Flink 算完推进 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 里等一个延迟窗口。窗口外的样本进离线补训。