实时数据平面
小红书把「用户和笔记互动 → 样本 → 模型 → 再推荐」写成实时闭环。能核对到的一线材料包括:阿里云收录的 Flink + Hologres 实践、Flink Forward Asia 2021(Native Flink on Kubernetes),以及 CNCF Volcano 官方博客里小红书推荐引擎那篇。
用户
App
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。场景包括实时反欺诈、实时数仓、实时推荐和数据传输。