在电商和导购类平台里,推荐系统的实时性几乎决定了用户的下单转化率。用户滑到某个商品卡片,系统必须在几百毫秒内判断“该不该推”“推哪个”,背后依赖的是一整套从实时特征计算到模型在线推理的链路。这篇文章我会完整拆解一套我参与设计的导购平台商品推荐引擎架构:基于Flink做实时特征工程,配合TensorFlow Serving承载在线推理,把用户从点击行为发生到推荐结果返回的延迟压到500毫秒以内。我会从整体设计思路、核心细节、实操过程、踩坑记录四个维度展开,尽量把每一步“为什么这么做”讲透,适合正在搭建或重构推荐系统的工程师参考。
1. 整体设计与思路拆解
1.1 为什么推荐引擎需要实时特征工程
传统的推荐系统大多走离线链路:凌晨用Spark批处理把前一天的用户行为、商品热度、类目偏好算好,写入特征库,白天在线服务直接读取。这种做法在商品动销率不高的平台够用,但放到导购平台场景下就有明显问题。
导购平台的特点是“人找货”和“货找人”并存,用户决策周期短,情绪化点击多。昨天用户看了某款口红但没买,今天早上另一款同色系口红突然因为某篇种草文火了,离线特征根本感知不到这个变化。用户的实时兴趣漂移、商品瞬时热度、当前场景的上下文信息,这些通通需要在秒级甚至毫秒级完成计算并反馈到推荐结果里。
我接手这套引擎时,产品方给的核心指标有两个:推荐结果点击率提升至少15%,用户从点击行为到刷出下一次推荐结果的响应时间不超过800毫秒。离线特征体系只能做到T+1更新,点击率天花板明显,响应时间倒是够,但效果有限。要想突破,只能上实时链路。
Flink在这里扮演的角色就是“实时特征计算引擎”,它负责把Kafka里源源不断的用户行为事件流、商品信息变更流、库存价格变更流做流式关联、窗口聚合、特征拼接,最后把计算好的特征实时写入在线存储供推理服务读取。整套链路里Flink不是唯一的组件,但它是所有实时特征计算的中枢。
1.2 Flink与TensorFlow Serving分工的本质边界
很多团队在做推荐系统时容易陷入一个误区:试图让Flink去做模型推理,或者让TensorFlow Serving顺带把特征拼接也干了。这两种做法我都见过,也都在生产环境里踩过坑。
Flink擅长的是无界流处理,它的状态管理、事件时间处理、窗口机制、checkpoint容错都是为“持续不断的数据流计算”设计的。但Flink不适合做高并发的单条请求推理,因为模型推理需要的是低延迟、高吞吐的服务能力,这正好是TensorFlow Serving这类专用推理服务器的强项。
TensorFlow Serving的优势在于:模型版本管理、热加载、自动批处理(dynamic batching)、基于gRPC的高效通信。它只做一件事——接收特征向量,返回预测分数。但它不负责特征计算,你把原始行为日志直接塞给TensorFlow Serving没有任何意义,模型要的是拼接好、清洗好、规整成固定维度的特征向量。
所以整个架构的本质分工是:Flink负责“从原始事件到可用特征”的流式加工,TensorFlow Serving负责“从特征向量到预测分数”的高性能计算。中间的桥梁是一个在线特征存储,通常用Redis或者阿里的Tair这类kv存储,Flink算好的特征写进去,推理服务启动时读出来拼装成完整向量。
1.3 方案选型背后的对比与取舍
在定这个架构之前,我对比过几套替代方案,这里把选型逻辑记录下来,方便后来人参考。
第一套备选方案是“纯Spark Streaming + Redis + PMML在线推理”。Spark Streaming的微批模式延迟在秒级,对于导购场景的实时性要求来说偏慢,而且Spark Streaming的本质是批处理,状态管理和窗口计算不如Flink灵活。PMML虽然能做到跨平台模型部署,但只支持逻辑回归、GBDT这类传统模型,深度学习模型完全没法走这条路。
第二套备选方案是“Flink + Flink ML + 自研推理服务”。Flink ML在做在线学习方面确实有潜力,但当时它的生态成熟度还不足以支撑生产级在线推理,尤其是深度学习模型这块,Flink ML的算子覆盖远不如TensorFlow完整。
第三套就是最终采用的“Flink + TensorFlow Serving”。这套组合的优势在于:Flink的实时计算能力和TensorFlow Serving的专用推理能力都是各自领域经过大规模验证的,两者之间通过特征存储解耦,各自可以独立扩容和升级。缺点是需要维护两个分布式系统,运维成本更高,但这个代价换来的实时性和灵活性在当时看来是值得的。
2. 实时特征工程的核心细节与实操要点
2.1 事件接入层:Kafka Topic的设计与序列化选型
实时特征工程的起点是行为事件接入。我们当时在Kafka里按事件类型建了多个Topic:user_click(用户点击)、user_view(商品曝光)、user_collect(收藏)、user_cart(加购)、order_create(下单)、item_info_change(商品信息变更)、price_change(价格变动)、stock_change(库存变动)。
Topic分得细有几个好处:一是下游Flink作业可以按需订阅,不需要消费全量数据再过滤;二是不同事件的吞吐量差异很大,点击数每秒上万条,下单数每秒可能只有几十条,分开建Topic方便分别调整分区数和消费并发。
序列化方案我们用了Avro,配合Confluent Schema Registry做Schema管理。当时很多团队直接用JSON,我看过他们的生产环境,最大的痛点是字段变更时下游作业全部要跟着改,稍有不慎就出现反序列化异常。用Avro之后,Schema演进方便很多,字段新增和废弃都有明确的兼容性规则。
这里有一个实操细节很多人忽略:Flink消费Kafka数据时,如果使用Avro序列化,建议在Flink作业里指定专用的DeserializationSchema,而不是用Flink自带的GenericRecord。直接转成POJO对象可以让后续的算子逻辑清晰很多,也能利用Flink的类型系统在序列化时拿到更优的性能。
2.2 窗口计算与状态管理:用户实时兴趣画像的构建
推荐引擎里最重要的实时特征就是用户当前的兴趣向量。我采用的方案是Flink的滑动窗口加KeyedState实现。
具体做法:以userId为key,开一个10分钟长度、5秒滑动一次的窗口,统计用户在这段时间内点击过的商品类目分布、价格带分布、浏览过的店铺List。窗口计算用Flink的增量聚合函数,每次只保留聚合结果,不缓存窗口内的所有明细数据,避免状态量过大。
窗口计算得出的类目分布是一个稀疏向量,比如“美妆0.6、服饰0.2、数码0.2”,这个向量会直接写入特征存储的user_realtime_profile字段。推荐服务读取这个字段时,就能知道这个用户此刻的兴趣重点在哪个类目上。
但滑动窗口有一个问题:窗口长度固定,用户兴趣变化的速度不一致。有人逛了5分钟就下单,有人逛了一小时还在犹豫。我当时用了一个折中的办法,同时维护多个不同时间长度的窗口特征:3分钟短期兴趣、15分钟中期兴趣、60分钟长期兴趣,三个维度的特征拼接在一起,让模型自己学习不同场景下应该更相信哪个时间尺度的信号。
状态管理方面,Flink的KeyedState默认使用RocksDB作为状态后端时会有序列化和反序列化开销,但对大状态场景几乎是唯一选择。我的建议是:如果每个用户的状态只有几KB,且并行度不高,用FsStateBackend配合Heap就可以;如果用户量大、状态量大,就必须上RocksDB,同时把state.backend.rocksdb.memory.managed配置调大,避免频繁的磁盘IO。
2.3 维表关联:实时拼接商品画像的关键路径
实时特征不能只算用户侧,商品侧的实时属性同样关键。商品标题、类目、品牌、价格、库存这些信息变化频率不高,但会在关键时刻影响推荐结果——比如某个商品突然降价了,它应该立即获得更高的推荐权重。
Flink关联维表有几种常见方案:查外部存储(Redis或MySQL)、广播维表、异步IO。我最终用的组合是:高频维表走广播流,低频维表走异步IO查Redis。
广播维表的适用场景是维表数据量不大(几万条以内)、更新频率不高、但每条数据都需要跟主数据流做关联。我们的商品基础信息维表虽然总量有几十万条,但热门商品的子集大概两万条,符合广播的条件。将这份热点商品维表做成BroadcastStream,每个并行子任务都保存一份全量数据,关联时纯内存查找,速度快且不依赖外部存储。
低频维表走异步IO,这里的“低频”指的是关联操作不是每条事件都触发,比如用户在浏览一件衣服时,我需要顺带查一下该店铺的30天销量和评分,这个数据放在Redis里,用Flink的AsyncFunction异步发起查询,配合连接池和超时控制,可以显著提升吞吐。
实操时有一个坑不得不提:维表关联的时间对齐问题。广播维表是定期更新的,更新前后同一个key可能关联到不同的值。如果关联的特征用于模型打分,这个不一致性会导致特征跳变。解决方案是给维表数据加一个version时间戳,写入特征存储时记录版本号,服务端读取特征时如果发现同一批特征里版本不一致,以多数版本为准或者丢弃该次推荐。
2.4 特征存储选型与特征拼接的细节
特征存储是整个实时特征工程和在线推理的中转站。我们用的是Redis Cluster,主要看中它的读写延迟低、集群模式支持水平扩展。
Redis中key的设计遵循一个原则:特征按不同的粒度拆分,避免大key。用户粒度特征key是rec:user:{userId}:profile,商品粒度特征key是rec:item:{itemId}:profile,上下文特征key是rec:ctx:{scene}:{date}:{hour}。每个key的value都是一个Hash,字段名是特征名,字段值是特征值。
Hash结构的优点是:推荐服务读取特征时可以只取自己关心的字段,不需要反序列化整个对象。缺点是字段多了之后内存开销大。所以特征数量控制在每个key不超过50个字段,超过的部分拆到第二个key。
特征拼接的细节在推理侧,但特征存储侧已经决定了拼接的效率。我当时的做法是让Flink把同一用户的所有实时特征聚合到一个Redis key里,推理服务一次MGET就能拿到所有在线特征,再加上离线服务预取的静态特征,拼成完整向量只需要几毫秒。
关于特征拼接还有一个容易忽略的坑:特征字段的顺序。TensorFlow Serving的模型输入是固定的FeatureSpec,训练时定义了什么顺序,线上推理就必须一模一样。所以Flink写特征存Redis时要按模型特征清单的字段顺序来,不能随意调整,否则推理结果会打折扣。
3. 在线推理架构的搭建与核心环节实现
3.1 TensorFlow Serving的部署与模型管理
TensorFlow Serving的部署本身不复杂,我踩过的坑都在模型管理、版本控制和生产调优上。
模型训练完成后导出为SavedModel格式,这是TensorFlow Serving唯一认识的模型格式。导出时需要指定serving_input_receiver_fn,它定义了线上推理时输入张量的结构。我建议在导出时就把特征向量的维度和dtype定死,不要用变长输入,变长输入会显著增加推理服务的调度开销。
我们当时用的是Docker部署,TensorFlow Serving的官方镜像直接拉下来,把模型目录挂载到容器里。模型目录里按版本号建子目录:models/rec_model/1、models/rec_model/2、models/rec_model/3,TensorFlow Serving启动时会自动扫描目录并加载最高版本的模型。
这里的版本控制有个细节:TensorFlow Serving默认加载最高版本模型,但新版本上线前需要灰度验证。通过配置model_config_list可以精确控制加载哪些版本,配合一个简单的流量切分逻辑——将小部分请求带上model_version=x的标签发到指定版本,验证通过后再切全量。
线上模型的更新频率也是一门学问。推荐模型我建议不要频繁更新,每天更新一次就够了,太频繁会导致线上效果波动,而且难以评估是模型更新贡献的还是实时特征贡献的。我们的节奏是凌晨离线训练,早上6点自动部署新模型,这样平台用户高峰期始终用的是当天最新的模型。
3.2 推理服务的请求-响应链路优化
用户请求进入推荐服务后,完整链路是:接收请求、解析用户上下文、从特征存储拉用户实时特征和商品特征、拼接成模型输入向量、通过gRPC调用TensorFlow Serving、拿到预测分数、按分数排序过滤、返回推荐列表。
这条链路每一步都可能成为瓶颈,我逐项优化过。gRPC调用的超时设置为50毫秒,超过直接降级,用离线分返回兜底,避免用户长时间等待。连接池大小设置为TensorFlow Serving最大并发数的1.5倍,充分利用其多线程模型。
TensorFlow Serving自带的动态批处理器(Dynamic Batching)值得好好研究。它会把一小段时间窗口内的多个请求合并成一个batch,一次性喂给模型推理,充分利用GPU或CPU的向量化能力。但批处理会增加第一个请求的等待时间,需要权衡。经过压测,我将batch_timeout_microseconds设置为3000微秒,即最多等3毫秒就发出批次,这样单请求延迟只增加几毫秒,但整体吞吐提升约30%。
3.3 GPU与CPU的选型取舍
TensorFlow Serving的推理硬件选型要看模型复杂度和QPS要求。我们的模型是DeepFM加上部分注意力结构,embedding维度不算大,单次推理耗时在CPU上约8毫秒,这个量级其实不需要上GPU。
GPU适合的场景是单次推理耗时超过20毫秒、QPS要求极高、且模型结构复杂(比如包含大尺寸Transformer结构)。GPU的优势在于并行计算吞吐高,但推理延迟不一定比CPU低,因为引入了数据拷贝和kernel launch的开销。导购推荐场景的模型通常不会太重,用多核CPU实例配合动态批处理反而性价比更高。
我们最终用的是8核16G的容器实例,每个实例部署一个TensorFlow Serving,显式设置了CPU线程数为4,并发数上限为200。容量规划是按照高峰期QPS的2倍冗余配置,确保大促时不会被打垮。
3.4 模型全链路一致性验证
这部分是最容易被忽略的。特征工程在Flink侧计算,模型在TensorFlow侧训练,两边各干各的,如果特征计算逻辑和训练时的特征处理逻辑不一致,再好的模型也是白搭。
我要求算法团队在训练时把特征处理的代码流程完整记录下来,包括缺失值填充方式、归一化参数(均值和方差)、类别特征的映射字典。这些元信息统一存放在配置中心里,Flink作业和推理服务都从配置中心读取同一份配置。
上线前必须做一致性校验:取当天的真实日志数据,离线批量跑一遍特征工程,得到离线特征;同时把日志同速回放到Flink作业里,得到在线特征;然后比较两者差异。误差允许在万分之一以内,超过这个阈值就需要排查原因。
4. 常见问题与排查技巧实录
4.1 Flink作业停止消费或延迟飙高的排查
Flink作业最常见的故障是消费延迟持续上涨,这个问题的排查思路我总结成一个固定流程。
先看Kafka的消费组Lag指标,确认到底是Flink没拉数据还是拉了数据但处理不过来。Flink这边看Web UI的Backpressure状态和繁忙程度,如果subtask显示HIGH状态,说明处理能力到瓶颈了;再看GC日志,如果频繁Full GC,多半是状态后端有问题。
状态后端的坑我遇到过一次典型场景:大促期间用户量暴增,每个人维护的窗口状态数据量翻了好几倍,RocksDB的写放大导致磁盘IO成瓶颈。当时紧急调整了并行度,从32扩到48,同时把Sink端的写入批量加大,延迟从叶秒级降回了秒级以内。
还有一个隐蔽的坑:Checkpoint超时导致作业反复重启。Checkpoint interval设置了60秒,但超时时间默认只有10分钟。大促时状态量大,一次checkpoint可能要8分钟,勉强能过,后来加了一次维度表和状态的数据增大,checkpoint直接超时,作业不断把状态回滚到上一个checkpoint,形成恶性循环。后来把timeout调整成20分钟,同时优化了RocksDB的压缩策略,才算稳定。
4.2 TensorFlow Serving内存泄漏与OOM
TensorFlow Serving本身比较稳定,但长时间运行后大概率会遇到内存持续上涨的问题。排查时我先用jstat观察堆内存变化,发现Old区占用量只增不减,GC后无法回收。
后来定位到是动态批处理模式下请求张量的生命周期过长。查了官方文档,确认是默认配置下请求队列的容量过大,导致积压的请求张量占用了大量堆内存。调整方案是显式设置max_batch_size和batch_timeout_microseconds,同时把模型加载数量从多个版本改为只保留最新两个版本。
这个问题的隐藏风险在于:内存泄漏不会立刻引发OOM,而是逐渐蚕食堆内存,最终在高峰期触发OOM导致容器重启。容器的自动重启机制虽然能快速恢复,但这段时间的请求全部走了降级策略,推荐质量大幅下滑,直接影响当天的转化率指标。所以对推理服务的监控要特别关注老年代内存的增长趋势,设置合理阈值提前预警。
4.3 特征不一致引发的线上效果波动
这应该是推荐系统最让人头疼的问题了。现象是模型刚上线那两天效果很好,第三天开始点击率逐日下滑,无论是回滚模型版本还是重新部署,效果都回不到上线初期的水平。
排查了几天才找到元凶:Flink作业在凌晨重启时,消费Kafka的offset没有对齐,导致部分用户的行为事件没有被计算进特征。这些用户的实时兴趣画像还停留在重启前的状态,而模型的输入分布变了,打分结果自然就偏了。
这类问题防不胜防,我后来干脆写了一版特征自检工具:在Redis里对每个用户存一个特征最后更新时间字段,如果发现某个用户的特征停留时间异常地长,就自动触发一次该用户的特征重算。另外在Flink作业重启流程里强制加入offset校准步骤,先暂停消费、确认checkpoint完成、再启动新作业,从代码层面杜绝这类隐患。
4.4 常见问题速查表
| 问题现象 | 排查方向 | 解决方案 |
|---|---|---|
| Flink消费延迟持续上涨 | 查看subtask繁忙度、GC日志、RocksDB IO | 调整并行度,优化状态后端配置,增大批量 |
| Checkpoint反复超时失败 | 查看checkpoint大小和耗时曲线 | 调大超时时间,优化增量checkpoint策略 |
| 推理响应时间突然变长 | 看TensorFlow Serving的CPU和内存瓶颈 | 调整动态批处理参数,扩容推理实例 |
| 推荐结果明显偏离用户近期行为 | 检查实时特征是否长时间未更新 | 自检工具主动重算,修复offset对齐问题 |
| 模型更新后效果剧烈下降 | 对比新旧版本的特征分布和评分分布 | 灰度切流,快速回滚,排查特征一致性问题 |
| Redis连接数打满 | 检查Flink异步IO连接池和推理服务连接池大小 | 调整连接池上限,增加Redis集群分片数 |
5. 一些经验心得
整套系统从设计到上线耗时两个多月,中间踩了很多坑,最后稳定运行后的架构比最初的设想简化了不少。一个很深的体会是:实时推荐系统的难点不在于单个技术组件有多强,而在于所有组件拼接处的一致性。
Flink侧保证的是“特征算得准”,TensorFlow Serving保证的是“分数推得对”,但两端一旦出现口径不一致,整个系统就会出现各种玄学故障。所以我在团队里立了一个规矩:每个特征从定义、计算、存储到模型消费,必须有一个唯一owner来负责端到端的正确性,任何一环的改动都要全链路回归验证。
另外关于团队配置,一个能跑通且稳定的实时推荐系统至少需要:一个精通Flink流式计算的数据工程师、一个熟悉TensorFlow Serving部署和模型调优的算法工程化工程师、一个能写高质量在线服务的后端工程师。这三个角色缺一不可,指望其中任何一个人包揽全部都不现实。
最后说一个可以继续扩展的方向:目前这套架构里Flink只负责了特征计算,事实上Flink还可以承担实时样本生成和在线学习的职责,把用户实时的反馈闭环到模型的增量更新中。这块我当时因为精力原因没有深入,但架构上是完全兼容的,Flink产出的实时样本可以作为在线学习的输入,训练好的增量模型再通过TensorFlow Serving的热加载机制无缝切换,这样整个推荐系统就拥有了真正意义上的自我进化能力。有条件的团队建议往这个方向走,实时推荐系统的最终形态一定是“感知-计算-推理-学习”一体化的。