推荐模型在线学习特征实时更新比模型在线更新更实用一、在线学习这个词被滥用了模型更新和特征更新是两回事谈到推荐模型的时效性在线学习几乎是每个技术分享的高频词汇。很多团队的想法是用户点了什么模型立刻吸收下一次推荐就变得更精准。这个愿景很美好但工程现实很残酷。真正的模型在线更新Online Learning意味着模型的权重在持续变化。TensorFlow 或 PyTorch 训练的模型不是一段静态的推理图而是一个不断吸收新样本、不断更新梯度的动态系统。这在广告竞价场景如 CTR 预估中确实有应用——因为广告点击信号极强、反馈极快延迟 5 分钟更新模型权重就能多赚 3% 的营收。但在推荐场景下在线更新模型面临三个现实制约。第一训练数据质量不稳定。用户点击不一定是正反馈——误触、快速划过、看了两秒就关这些噪声如果被实时吸收进模型权重会被污染。第二模型在线更新的工程复杂度极高。需要实现参数服务器Parameter Server保证分布式梯度同步的一致性处理 Worker 节点故障时的梯度回滚。第三推理和训练共享 GPU 资源线上流量高峰期训练任务会抢推理任务的 GPU导致推荐延迟毛刺。相比之下特征实时更新是一个投入产出比高得多的方案。模型保持相对稳定的权重每天更新一次但输入给模型的特征向量是实时变化的。用户最近 5 分钟点了什么、浏览了什么品类、停留时长多少这些行为不改变模型的参数但改变了模型看到的输入画面。效果上特征实时更新可以覆盖在线学习 80% 以上的收益工程复杂度却只有后者的 20%。二、特征实时更新的管道设计Kafka → Flink → Redis特征实时更新的架构是一条标准的三段式流处理管道。数据采集端所有用户行为通过埋点 SDK 或服务端日志进入 Kafka。Topic 设计有讲究不同行为类型分不同 Topic如user_click、user_view、user_purchase因为不同行为的权重不同——点击权重 1购买权重 5浏览权重 0.3。合并到一个 Topic 后续计算要额外增加过滤逻辑。实时计算端Flink 从 Kafka 消费事件按用户 ID 做 KeyBy 分组在时间窗口内聚合特征。窗口选择是工程决策的关键滚动窗口Tumbling Window5 分钟计算最近 5 分钟的品类偏好分布。窗口边界对齐到分钟级整点无重叠事件不重复计数。滑动窗口Sliding Window30 分钟滑动步长 5 分钟计算最近 30 分钟的点击频次。窗口之间有大量重叠同一事件在多个窗口中计数——这不是 bug这是统计窗口的预期行为。会话窗口Session Window30 分钟间隔聚合一次使用会话内的行为序列。两次行为间隔超过 30 分钟视为不同会话适合计算单次会话内的兴趣漂移。Flink 输出端直接写 Redis。为了减少 Redis 写入压力Flink 算子内聚合并批写入——攒够 100 条或超过 1 秒才 flush 一次。单条写入 Redis 的延迟约 1ms批量写入 100 条的延迟约 2-3ms吞吐量提升 30 倍以上。特征消费端Go 特征服务从 Redis 读取这些实时特征。为了降低 Redis 的读取压力特征服务内部维护了一层本地 LRU 缓存TTL 5 分钟——和 Flink 窗口周期对齐。只要窗口没切换同样的用户从本地缓存就能拿到特征数据。命中率实测约为 75%节省了对应比例的 Redis 查询。三、模型日更新的流程离线训练不会消失特征实时更新承担了时效性的任务但模型更新也不能完全废弃。模型需要定期用全量样本重新训练这个周期通常是一天一次。为什么不能只用特征更新、永远不更新模型因为特征分布会漂移。用户行为模式在变化——夏天推短袖冬天推羽绒服——如果模型永远用三个月前的数据训练即使特征是最新的气温数据模型也不知道低温和羽绒服之间的关联权重该是多少。模型参数需要周期性地吸收新样本让权重跟上数据分布的变化。日更新的流程是标准化的离线训练链路凌晨 2 点Spark 任务从 HDFS 拉取过去 30 天的全量行为日志做样本拼接正样本点击负样本曝光未点击凌晨 3 点样本按照 8:1:1 切分训练集、验证集、测试集凌晨 4-6 点PyTorch/TensorFlow 分布式训练完成后将模型导出为 SavedModel 格式早上 7 点新模型推送到模型仓库OSS/S3推理服务通过 Triton 的 Model Repository 自动热加载模型日更新 特征秒级更新这个组合在工业界被验证是最平衡的时效性方案。模型参数是慢变量决定了推荐系统的全局品味特征数据是快变量决定了推荐系统对用户当下意图的感知速度。四、特征实时更新的边界数据倾斜和状态爆炸特征实时更新的方案不是全无代价的。Flink 处理用户行为流时会遇到两类典型的流计算问题。数据倾斜。少数头部用户的行为量是普通用户的百倍以上——一个重度用户每分钟产生几十次点击Flink 算子里这个用户的 Key 对应的 Task 处理数据量远大于其他 Task。如果不做处理这个 Task 的处理延迟会飙高拖慢整个窗口的输出。解决方案是两级聚合第一级在 Source 侧做局部聚合每个 TaskManager 内部先聚合自己负责的数据第二级再做全局 KeyBy 聚合。这样可以有效摊薄热点 Key 的负载。状态膨胀。Flink 的有状态计算意味着每个用户都要在内存中维护一个窗口内的行为累积状态。日活 100 万的系统每个用户对应一组聚合状态约 1KB累计 1GB。加上 Flink 的 Checkpoint 机制需要将状态快照写入 RocksDB——状态越大Checkpoint 越慢。配置 Checkpoint 间隔为 5 分钟时1GB 状态写入 RocksDB 约需 10-15 秒。如果状态膨胀到 10GBCheckpoint 需要 1 分钟以上延迟不可接受。控制状态大小的关键策略是设置合理的状态 TTL非活跃用户30 分钟无新事件的状态自动清理减少 60% 以上的无效状态存储。还有一个现实问题是数据延迟 vs 数据精度的取舍。Flink 使用 Event Time 处理允许一定程度的乱序数据watermark 延迟 30 秒。这意味着用户在 10:00:01 秒的点击可能在 10:00:31 秒才进入窗口计算。对于推荐场景这个延迟可以接受——用户在 10:00 看到的推荐用的是 09:59:30 秒的特征30 秒的滞后在体验上几乎无感知。五、总结推荐系统时效性的正确打开方式不是我全都要做在线学习而是特征实时更新 模型日更新的分工协作。特征管快——秒级反映用户当下意图通过 Kafka → Flink → Redis 管道实现模型管稳——天级更新全局参数通过离线 Spark 训练 Triton 热加载实现。这套方案在工程实现中需要关注 Flink 的数据倾斜处理、状态 TTL 管理和 Event Time 容错。特征层的 Redis 缓存加上应用层的本地缓存可以形成两层防护——即使 Flink 偶尔延迟推荐服务也不会立刻降级。场景适配的底线搜索广告、实时竞价等极低延迟反馈的场景在线学习仍然有价值。但在 90% 的内容和电商推荐场景下特征实时更新提供了足够的效果和可维护性——不需要为了追求在线学习这个听起来很酷的词去支付参数服务器和分布式梯度同步那笔巨额的工程账单。
推荐模型在线学习:特征实时更新比模型在线更新更实用
推荐模型在线学习特征实时更新比模型在线更新更实用一、在线学习这个词被滥用了模型更新和特征更新是两回事谈到推荐模型的时效性在线学习几乎是每个技术分享的高频词汇。很多团队的想法是用户点了什么模型立刻吸收下一次推荐就变得更精准。这个愿景很美好但工程现实很残酷。真正的模型在线更新Online Learning意味着模型的权重在持续变化。TensorFlow 或 PyTorch 训练的模型不是一段静态的推理图而是一个不断吸收新样本、不断更新梯度的动态系统。这在广告竞价场景如 CTR 预估中确实有应用——因为广告点击信号极强、反馈极快延迟 5 分钟更新模型权重就能多赚 3% 的营收。但在推荐场景下在线更新模型面临三个现实制约。第一训练数据质量不稳定。用户点击不一定是正反馈——误触、快速划过、看了两秒就关这些噪声如果被实时吸收进模型权重会被污染。第二模型在线更新的工程复杂度极高。需要实现参数服务器Parameter Server保证分布式梯度同步的一致性处理 Worker 节点故障时的梯度回滚。第三推理和训练共享 GPU 资源线上流量高峰期训练任务会抢推理任务的 GPU导致推荐延迟毛刺。相比之下特征实时更新是一个投入产出比高得多的方案。模型保持相对稳定的权重每天更新一次但输入给模型的特征向量是实时变化的。用户最近 5 分钟点了什么、浏览了什么品类、停留时长多少这些行为不改变模型的参数但改变了模型看到的输入画面。效果上特征实时更新可以覆盖在线学习 80% 以上的收益工程复杂度却只有后者的 20%。二、特征实时更新的管道设计Kafka → Flink → Redis特征实时更新的架构是一条标准的三段式流处理管道。数据采集端所有用户行为通过埋点 SDK 或服务端日志进入 Kafka。Topic 设计有讲究不同行为类型分不同 Topic如user_click、user_view、user_purchase因为不同行为的权重不同——点击权重 1购买权重 5浏览权重 0.3。合并到一个 Topic 后续计算要额外增加过滤逻辑。实时计算端Flink 从 Kafka 消费事件按用户 ID 做 KeyBy 分组在时间窗口内聚合特征。窗口选择是工程决策的关键滚动窗口Tumbling Window5 分钟计算最近 5 分钟的品类偏好分布。窗口边界对齐到分钟级整点无重叠事件不重复计数。滑动窗口Sliding Window30 分钟滑动步长 5 分钟计算最近 30 分钟的点击频次。窗口之间有大量重叠同一事件在多个窗口中计数——这不是 bug这是统计窗口的预期行为。会话窗口Session Window30 分钟间隔聚合一次使用会话内的行为序列。两次行为间隔超过 30 分钟视为不同会话适合计算单次会话内的兴趣漂移。Flink 输出端直接写 Redis。为了减少 Redis 写入压力Flink 算子内聚合并批写入——攒够 100 条或超过 1 秒才 flush 一次。单条写入 Redis 的延迟约 1ms批量写入 100 条的延迟约 2-3ms吞吐量提升 30 倍以上。特征消费端Go 特征服务从 Redis 读取这些实时特征。为了降低 Redis 的读取压力特征服务内部维护了一层本地 LRU 缓存TTL 5 分钟——和 Flink 窗口周期对齐。只要窗口没切换同样的用户从本地缓存就能拿到特征数据。命中率实测约为 75%节省了对应比例的 Redis 查询。三、模型日更新的流程离线训练不会消失特征实时更新承担了时效性的任务但模型更新也不能完全废弃。模型需要定期用全量样本重新训练这个周期通常是一天一次。为什么不能只用特征更新、永远不更新模型因为特征分布会漂移。用户行为模式在变化——夏天推短袖冬天推羽绒服——如果模型永远用三个月前的数据训练即使特征是最新的气温数据模型也不知道低温和羽绒服之间的关联权重该是多少。模型参数需要周期性地吸收新样本让权重跟上数据分布的变化。日更新的流程是标准化的离线训练链路凌晨 2 点Spark 任务从 HDFS 拉取过去 30 天的全量行为日志做样本拼接正样本点击负样本曝光未点击凌晨 3 点样本按照 8:1:1 切分训练集、验证集、测试集凌晨 4-6 点PyTorch/TensorFlow 分布式训练完成后将模型导出为 SavedModel 格式早上 7 点新模型推送到模型仓库OSS/S3推理服务通过 Triton 的 Model Repository 自动热加载模型日更新 特征秒级更新这个组合在工业界被验证是最平衡的时效性方案。模型参数是慢变量决定了推荐系统的全局品味特征数据是快变量决定了推荐系统对用户当下意图的感知速度。四、特征实时更新的边界数据倾斜和状态爆炸特征实时更新的方案不是全无代价的。Flink 处理用户行为流时会遇到两类典型的流计算问题。数据倾斜。少数头部用户的行为量是普通用户的百倍以上——一个重度用户每分钟产生几十次点击Flink 算子里这个用户的 Key 对应的 Task 处理数据量远大于其他 Task。如果不做处理这个 Task 的处理延迟会飙高拖慢整个窗口的输出。解决方案是两级聚合第一级在 Source 侧做局部聚合每个 TaskManager 内部先聚合自己负责的数据第二级再做全局 KeyBy 聚合。这样可以有效摊薄热点 Key 的负载。状态膨胀。Flink 的有状态计算意味着每个用户都要在内存中维护一个窗口内的行为累积状态。日活 100 万的系统每个用户对应一组聚合状态约 1KB累计 1GB。加上 Flink 的 Checkpoint 机制需要将状态快照写入 RocksDB——状态越大Checkpoint 越慢。配置 Checkpoint 间隔为 5 分钟时1GB 状态写入 RocksDB 约需 10-15 秒。如果状态膨胀到 10GBCheckpoint 需要 1 分钟以上延迟不可接受。控制状态大小的关键策略是设置合理的状态 TTL非活跃用户30 分钟无新事件的状态自动清理减少 60% 以上的无效状态存储。还有一个现实问题是数据延迟 vs 数据精度的取舍。Flink 使用 Event Time 处理允许一定程度的乱序数据watermark 延迟 30 秒。这意味着用户在 10:00:01 秒的点击可能在 10:00:31 秒才进入窗口计算。对于推荐场景这个延迟可以接受——用户在 10:00 看到的推荐用的是 09:59:30 秒的特征30 秒的滞后在体验上几乎无感知。五、总结推荐系统时效性的正确打开方式不是我全都要做在线学习而是特征实时更新 模型日更新的分工协作。特征管快——秒级反映用户当下意图通过 Kafka → Flink → Redis 管道实现模型管稳——天级更新全局参数通过离线 Spark 训练 Triton 热加载实现。这套方案在工程实现中需要关注 Flink 的数据倾斜处理、状态 TTL 管理和 Event Time 容错。特征层的 Redis 缓存加上应用层的本地缓存可以形成两层防护——即使 Flink 偶尔延迟推荐服务也不会立刻降级。场景适配的底线搜索广告、实时竞价等极低延迟反馈的场景在线学习仍然有价值。但在 90% 的内容和电商推荐场景下特征实时更新提供了足够的效果和可维护性——不需要为了追求在线学习这个听起来很酷的词去支付参数服务器和分布式梯度同步那笔巨额的工程账单。