learntoall

实时数据平面

小红书把「用户和笔记互动 → 样本 → 模型 → 再推荐」写成实时闭环。能核对到的一线材料包括:阿里云收录的 Flink + Hologres 实践、Flink Forward Asia 2021(Native Flink on Kubernetes),以及 CNCF Volcano 官方博客里小红书推荐引擎那篇。

用户
App
Kafka
Flink
FeatureJoiner
模型

点击 / 收藏 / 隐藏

打点

归因和标签汇总

和笔记特征拼接

训练样本

下次发现页用上新模型
时序图 · 打点进 Kafka,Flink 归因,FeatureJoiner 出样本

推荐训练链路

步骤公开说法
采集前后端打点进入 Kafka(或 RocketMQ,视分享年代)。
归因Flink 根据打点生成行为标签;有效点击等复杂规则希望只实现一次。
汇总再一个 Flink 作业做出 Summary 标签。
FeatureJoiner标签和推荐引擎里的笔记特征在 Flink 关联,得到训练样本。
训练样本进模型;Firefly 调度 TensorFlow / LarC。
分析实时指标进 ClickHouse 或 Hologres;离线样本落 Hive。

批流同一套逻辑

小红书强调:广告或推荐里「点击后停留超过数秒才算有效」这种规则,如果在 Flink 和离线 SQL 各写一遍,一定会出现两个有效点击。他们用 Flink 的批流一体(分享中提到 FLIP-27)让同一份代码既能吃日志文件,也能吃流。

百川、Firefly、Volcano

  • 百川(Baichuan):流计算平台,管实时标签和在线学习相关的 Flink 作业。
  • Firefly:机器学习任务平台。
  • LarC:基于 TensorFlow 的搜推广稀疏大模型训练框架。
  • Volcano:Kubernetes 上的批调度,用来更新实时和批量模型。

多云

Flink Forward Asia 2021:业务数据分散在阿里云、腾讯云、华为云,Flink 集群跟着走,checkpoint 落在对应的 OSS / COS / OBS。当时内部 SQL 与 JAR 任务比约 9:1。场景包括实时反欺诈、实时数仓、实时推荐和数据传输。