关于AI在线特征服务的几点记录
先复现,再谈优化。
**
我们需要一个在线特征服务,核心诉求就这么几条:
问题背景
离线特征批处理的时候觉得还挺顺手,每天跑一次 Spark Job,把用户行为、商品属性、历史统计数据全算好,存到 Hive 里。模型训练时从 Hive 读数据,推理时把特征预计算好放到 MySQL,看起来一切都挺美好的。
直到上线后才发现问题来了:
- 新用户注册后,离线特征还没算出来,推荐系统只能给个冷启动兜底策略
- 用户刚点了几个商品,兴趣特征还是昨天的数据,推荐列表完全跟不上节奏
- 运营上了一个新活动,特征更新要等到第二天凌晨才能生效
- 热点商品突然爆火,但热度特征还是昨天的冷数据
问题很清楚:离线批处理解决不了"实时"这件事。
需求梳理
我们需要一个在线特征服务,核心诉求就这么几条:
- 低延迟访问:推理请求要能在一两百毫秒内拿到特征值,不能等太久
- 准实时更新:用户行为发生后,特征最好能在几秒到几分钟内更新
- 高可用:特征服务挂了,推荐就彻底瘫痪,不能有单点故障
- 易扩展:随着数据量和 QPS 上涨,系统要能横向扩展
还要考虑现实限制:
- 团队资源有限,不能从零开始造轮子,要能快速落地
- 现有的离线特征pipeline不能丢,要能平滑迁移
- 成本要控制,不能为个特征服务上整套大数据平台
实现方案
架构设计
最终采用的架构大概是这个样子:
这个架构的核心思路是:
- 用户行为事件写入 Kafka:前端、后端、移动端的行为数据都统一发到 Kafka,作为数据源
- 实时特征计算:消费 Kafka 消息,计算需要实时更新的特征
- Redis Cluster 存储:计算好的特征直接写入 Redis,提供低延迟查询
- 离线特征补全:对于那些不需要实时、但查询时又必须有的特征,从离线批处理结果中补全
数据流转
整个数据流可以分为两条线:
实时线:用户行为 → Kafka → Flink/Spark Streaming → Redis 离线线:历史数据 → Spark Batch → Hive → Redis(通过离线特征补全服务)
实时线负责那些时间敏感的特征,比如最近N分钟的点击数、当前会话的行为偏好等。离线线负责那些计算成本高、不要求秒级更新的特征,比如用户LTV预测、商品长期热度趋势等。
技术选型
消息队列:Kafka
- 用 Kafka 而不用别的,主要是团队有经验,运维成熟
- 分区机制天然支持横向扩展,QPS 上来后加分区就行
- 消息持久化,特征计算服务挂了也不丢数据
计算引擎:Flink
- 最早想过用 Spark Streaming,但 Flink 的延迟确实更低
- 窗口机制、状态管理都比较完善,写实时特征计算很顺手
- 社区活跃,文档齐全,踩坑能找到人讨论
存储:Redis Cluster
- 单机 Redis 不够用,QPS 上去后扛不住
- Cluster 模式支持横向扩展,数据自动分片
- 内存存储,读取延迟确实能控制在毫秒级
离线特征补全:Thrift 服务
- 用 Thrift 定义接口,多语言支持好
- 从 Hive 读离线特征,合并到 Redis 查询结果里
- 这个服务可以批量预加载,减少 Hive 查询压力
踩坑记录
坑一:Redis 内存爆炸
最开始把所有特征都往 Redis 里塞,没几天内存就爆了。问题出在:
- 特征版本没控制,每次更新都是新增 key,旧 key 没清理
- 用户基数太大,每个用户都存几百个特征字段
- 热点数据的访问频率太高,Redis 集群负载不均
后来做了几件事:
- 设置合理的过期时间:不同类型特征设不同 TTL,实时特征短一点,离线特征长一点
- 特征字段压缩:把频繁更新的字段和稳定字段分开,稳定字段用更紧凑的编码
- 热点数据分片:对那些访问特别频繁的特征,增加副本数,分散压力
坑二:特征更新延迟太大
理想情况是用户行为发生后几秒钟特征就能更新,但实际情况经常是几分钟甚至更久。
排查后发现几个问题:
- Kafka 消费积压:特征计算服务处理不过来,消息堆积在队列里
- 窗口设计不合理:Flink 的窗口太大,导致特征更新频率低
- 网络延迟:特征计算服务和 Redis 不在同一机房,跨机房访问有延迟
解决方案:
- 增加特征计算服务的实例数,提高消费能力
- 优化窗口设计,对时间敏感的特征用更小的窗口
- 把特征计算服务和 Redis 部署在同一机房,减少网络开销
坑三:特征一致性难以保证
离线特征和实时特征经常对不上,模型预测时用的特征和训练时用的特征不一样,导致效果变差。
主要问题:
- 特征定义不一致:离线计算和实时计算用的逻辑不一样
- 时间窗口对齐:离线用的是固定时间窗口,实时用的是滑动窗口
- 数据源差异:离线读的是 Hive 里的全量数据,实时只读到 Kafka 里的增量数据
应对措施:
- 统一特征定义,把特征计算逻辑抽象成通用函数,离线和实时都用同一套
- 严格对齐时间窗口,离线和实时都用相同的时间语义
- 做特征校验,定期对比离线和实时特征的计算结果,发现差异及时告警
结果与收益
折腾了几个月,总算把在线特征服务跑起来了。效果还算看得见:
- 推荐响应时间从 500ms 降到 150ms:主要是特征查询快了,整个推理链路都跟着提速
- 新用户冷启动效果提升 20%:实时特征能让推荐系统更快捕捉到用户兴趣
- 运营活动生效时间从 T+1 缩短到 T+5分钟:新上的活动特征能快速更新到推荐系统里
- 系统可用性达到 99.9%:Redis Cluster + 多副本部署,基本没出现过因为特征服务挂掉导致推荐不可用的情况
也不是没有代价:
- 运维复杂度明显上升,要同时维护离线批处理和在线实时两套系统
- 成本增加了不少,Redis Cluster、Kafka 集群都要算钱
- 团队学习成本提高,以前只会写 Spark SQL,现在还得懂 Flink、Redis 运维
还能做些什么
目前的在线特征服务只能说解决了"有无"的问题,还有很多优化空间:
- 特征预加载:把那些访问频率特别高的特征提前加载到本地缓存,进一步减少 Redis 查询
- 特征版本管理:做更完善的特征版本控制,支持 A/B 测试和灰度发布
- 特征监控告警:监控特征质量,发现异常数据及时告警
- 特征血缘管理:记录特征的计算路径,方便问题排查和特征优化
结语
从离线到实时,看起来只是个技术升级,但实际上是整个数据处理思维的转变。离线处理时可以追求"完整性"和"准确性",但在线服务得在"实时性"和"准确性"之间做取舍。
特征服务就像桥梁,连接着数据和应用。桥建得好不好,不仅看技术选型,更看对业务的理解和对限制条件的把握。很多时候,最优解不是最新的技术,而是最符合当前团队、当前业务、当前资源约束的那个方案。
折腾这事就是这样,看起来是技术问题,到最后都是权衡。
版权声明: 本文首发于 指尖魔法屋-关于AI在线特征服务的几点记录(https://blog.thinkmoon.cn/post/287-ai-online-feature-serving-offline-realtime-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!
评论
使用 GitHub 账号登录后即可留言,支持 Markdown。