最近在帮朋友优化一个导购平台的推荐链路,感触很深。很多团队做推荐,模型训练和特征工程都挺完善,但到了线上,用户从"点击商品"到"推荐位更新"之间隔着十几分钟甚至更久,转化率就被拖垮了。我们后来落地了一套基于Flink的实时特征工程 + TensorFlow Serving在线推理架构,把推荐引擎的时效性从分钟级拉到了秒级。这篇内容不是PPT架构图,是我们在真实业务里踩坑、调优、上线后的完整复盘,包括自定义Source和Sink的实战写法、JDBC连接器异常、CDC部署、Hive Sink入不了表这些具体问题,希望对做推荐和实时计算的同行有参考价值。
当时我们明确的目标很朴素:用户刚搜过"冬季连衣裙",回到首页推荐位必须能看到相关商品;用户加购了某品牌吹风机,详情页的"相似好物"要尽快跟着变。要做到这一点,光靠离线T+1的特征根本无法完成,必须在事件发生的瞬间把行为特征算出来,再扔给模型打分。
1. 导购推荐为什么从"离线打分"转向"实时特征+在线推理"
1.1 离线推荐模式暴露出的三个天花板
先说说我们之前用的离线推荐链路。每天凌晨,离线任务把用户历史行为、商品信息、类目偏好全部算好,产出用户特征和召回候选集,推到线上的KV存储里。白天用户请求过来,推荐服务从KV拿特征,跑一个比较轻的模型打分,把结果返回前端。
这套模式有两个显而易见的问题。第一是特征时延太长,用户当天产生的行为要第二天才能进特征,像"刚搜索完某个品类"这类强意图信号完全抓不住。第二是模型是"日级"更新的,白天用户意图已经变化了好几轮,模型还停留在昨天凌晨的状态。导购场景还有一个特殊痛点,就是商品和活动更新快,大促期间秒杀品、新品、限时券的生命周期可能只有几个小时,离线特征根本来不及反映。
我们做过一个统计,导购App里,用户在一次会话中的平均浏览时长大概是3-5分钟,但离线推荐特征的更新周期是24小时。也就是说,用户在这一段会话里产生的"加购""点赞""停留"信号,对当前推荐位完全不起作用,只能影响明天的推荐。当这个对比摆在面前,实时化的必要性就不需要再论证了。
1.2 实时特征工程和在线推理到底解决了什么
实时化之后,链路变成了这样:用户在App的行为通过埋点上报到消息队列,Flink实时消费,把行为做窗口聚合、行为序列编码、维度特征关联,秒级写入在线特征存储。推荐服务收到请求时,实时拼接特征,调用TensorFlow Serving完成打分。
这里有两个常被误解的点,我得说清楚。第一,"实时特征"不是简单把离线特征的计算频率从每天一次改成每秒一次,而是整个特征逻辑要重新设计。离线可以全量扫描用户所有历史,实时就只能靠流式状态和窗口去维护"最近N小时""最近N次行为"这类滑动视图。第二,TensorFlow Serving不是万能的,它只是一个高效的模型托管和推理服务,真正难的是特征组织、模型导出格式、降级策略。模型在离线训练时用了哪些特征、特征怎么预处理,在线打分时必须一模一样,这是最容易翻车的地方。
整体收益在AB实验里非常明显:首页推荐的点击率相对提升了18%左右,加购转化率提升了接近10%。更重要的是,新用户冷启动和爆款商品实时追热这两类场景,以前基本无解,现在靠实时行为序列能比较快地响应。
2. 整体架构:Flink实时特征管道与TensorFlow Serving推理层如何衔接
2.1 分层设计:从埋点到推荐结果的完整链路
我们的架构分成五层,每一层职责都比较单一,方便排查问题。第一层是数据接入层,App端和服务端的行为日志统一进Kafka,按行为类型分了多个Topic。第二层是实时特征层,Flink集群消费Kafka,做清洗、聚合、序列编码、维度join,产出实时特征。第三层是存储与索引层,Redis里存用户最近行为序列和实时特征向量,ES/HBase存商品画像和候选集。第四层是推理层,TensorFlow Serving加载SavedModel格式的模型,接收特征做打分。第五层是业务接入层,推荐服务写一个gRPC客户端,把请求打到Serving,拿回TopN结果。
链路大致是:埋点 -> Kafka -> Flink -> Redis/在线特征 -> 推荐服务 -> TensorFlow Serving -> TopN。这里最容易被忽略的是"推荐服务"这一层,它不只是转发请求,还要负责特征拼接、模型打分超时处理、无特征时的降级逻辑。我们曾经见过一个设计方案,推荐服务把所有数据拼好后直接请求Serving,完全没有降级,结果模型服务抖动一下,整个推荐位就空了,这就是架构上没想清楚。
2.2 各层之间的时延预算
做过实时系统的人都知道,"实时"如果不定义时延目标,后面验收一定会吵架。我们当时给每个环节定了明确的时延预算:
| 环节 | 目标时延 | 说明 |
|---|---|---|
| 埋点上报到Kafka | <= 1s | 客户端批量上报,网络抖动另算 |
| Kafka到Flink处理产出特征 | <= 2s | 视Topic分区数和并行度而定 |
| 特征写入Redis | <= 100ms | 批量Pipeline写入 |
| 推荐请求到Serving打分返回 | <= 50ms | 纯模型推理部分 |
| 端到端(用户行为发生到推荐可用) | <= 5s | 最关键的可观测指标 |
这里要强调,端到端时延和单点的处理时延是两回事。很多时候Flink处理很快,但埋点SDK因为网络策略延迟上报,或者Kafka Topic分区数不够导致消费积压,前端推荐位还是半天不更新。所以从第一天起,就必须在链路埋点里加上时间戳,用Flink的Watermark来度量从"行为发生"到"特征可用"的真实延迟。
2.3 Flink与TensorFlow Serving的版本选型
版本选型看起来很基础,但决定了很多坑的大小。Flink我们用的1.17版本,原因有三点:一是对Flink SQL的维表join支持更成熟,二是Kafka和JDBC连接器版本跟得比较紧,三是流式模式下状态TTL的语义更清晰。TensorFlow Serving用的2.x分支,模型导出直接用TensorFlow 2的SavedModel格式,通过tf.saved_model.save导出,而不是老的tf.contrib.predictor。
版本之间最要命的问题是序列化兼容。我们的模型输入用了tf.Example协议格式,特征名、特征类型都必须在导出时固定。如果训练代码里用了tf.train.Feature构建样本,在线请求也要拼出同样的tf.Example。这个一致性如果破坏,Serving端不会报错,但打分结果全是垃圾。我们在压测时就发现,REST接口传JSON和gRPC传tf.Example,同样的特征值打分结果居然有细微差别,后来查了半天,是JSON浮点数精度问题。所以线上统一用gRPC,这个教训后面细说。
3. Flink实时特征工程的关键代码与设计思路
3.1 用DataStream API还是Flink SQL
这是每个Flink项目都会遇到的灵魂拷问。我们从实战角度给出的建议是:以Flink SQL为主,以DataStream API为辅。SQL的好处是开发快、维表join语法成熟、血缘和优化器能帮忙做很多事情。但SQL不适合做需要精细控制状态的复杂事件处理,比如按用户维护一个最近100个行为的滑动窗口,SQL表达起来很别扭,DataStream的KeyedProcessFunction就清晰得多。
我们的分工是:数据清洗、简单聚合、维表关联用Flink SQL;行为序列编码、自定义触发逻辑、复杂状态管理用DataStream API。还有一个折中方案,如果团队SQL能力一般,可以全部用DataStream API加高阶函数,但一定要封装好,不要到处散落状态逻辑。
3.2 用户行为序列的实时编码:ProcessFunction + KeyedState
实时特征里最有价值的一部分,是把用户最近一段时间的连续行为变成一个向量或一组统计量,喂给模型。这个逻辑不适合用SQL窗口做,因为我们需要的是"任意时刻都能拿到用户最近200个行为的编码结果",而不是固定窗口结束才能拿到。我们实现了一个UserBehaviorSeqProcessFunction,本质上是KeyedState + 定时器。
核心思路是:按用户ID做keyBy,用ValueState<List<BehaviorItem>>保存最近N个行为,用MapState<String, Long>记录各行为类型的计数。每来一条新行为,先判断当前时间与上一条行为的间隔,如果超过30分钟,说明是新会话,要重置会话内状态。然后把行为压入列表,如果超过200条就剔除最旧的,并更新时间衰减因子。最后把编码后的特征向量发往下游。
这里有两个特别重要的细节。第一,状态大小必须收敛。如果无限保留用户所有历史行为,大促期间热用户的状态会暴涨,直接拖垮内存和Checkpoint。我们的做法是ListState只保留最近200条,每处理完一批就清理超出的部分,同时配合StateTtlConfig给状态设置24小时的TTL。第二,更新下游特征时要有版本号。Redis里同一个用户key可能被并发写入,需要记录一个事件时间戳,下游读取特征时能感知到数据新鲜度,而不是盲信任一个旧快照。
// 伪代码,展示核心状态管理逻辑 public class UserBehaviorSeqProcessFunction extends KeyedProcessFunction<String, BehaviorEvent, UserFeature> { private transient ListState<BehaviorItem> recentBehaviors; private transient ValueState<Long> lastEventTime; @Override public void processElement(BehaviorEvent event, Context ctx, Collector<UserFeature> out) throws Exception { Long lastTs = lastEventTime.value(); if (lastTs != null && event.ts - lastTs > SESSION_GAP_MS) { recentBehaviors.clear(); } recentBehaviors.add(new BehaviorItem(event.itemId, event.type, event.ts)); // 裁剪到最多200条 if (recentBehaviors.get().size() > MAX_SEQ_LEN) { recentBehaviors.update(recentBehaviors.get().subList(recentBehaviors.get().size() - MAX_SEQ_LEN, recentBehaviors.get().size())); } lastEventTime.update(event.ts); out.collect(buildUserFeature()); } }3.3 维度特征实时关联:Async I/O + 维表缓存
实时特征不只有行为序列,还要关联上商品信息、用户画像、价格区间等维度属性。如果每来一条行为都同步查一次MySQL,Flink的吞吐量会直接被打死。我们用了Async I/O + 本地缓存两层方案。
Async I/O是Flink官方提供的异步访问外部系统的API,核心优势是发一个异步请求后不阻塞算子线程,等回调返回后再继续处理。我们在AsyncFunction里并发查询用户画像和商品信息,查询结果拼接到特征上。但异步也只是解决了并发度问题,没解决热点key的重复查询问题。所以我们又加了一层Cache<Key, Value>,给维表查询结果设置30秒到5分钟不等的过期时间。热门商品和热门用户的维度信息缓存命中率很高,实测MySQL的QPS压力下降了80%以上。
维度关联最容易出错的点在"时区"和"默认值"。比如"今日是否大促"这种特征,业务侧的“今日”和Flink处理时的北京时间可能不一致,如果不统一用Asia/Shanghai时区计算,离线训练和在线推理会拿到完全不同的特征值。另外,关联不上的维度一定要有显式的默认值策略,不能返回null,否则模型推理时特征缺失,行为完全不可预期。
3.4 特征结果落库与双写策略
实时特征算出来之后,不能只写一个地方。我们的实践经验是双写:一份写到Redis供在线推理读取,TTL设置成6小时,避免Redis膨胀;另一份写回Kafka,供离线数仓消费,做特征回放和训练样本拼接。Kafka里的特征Topic会最终落到Hive表,用于验证"线上打分用的特征"和"离线训练用的特征"是否一致,这个一致性验证非常重要。
写Redis的方式也经历过一次优化。早期是一条一条SET,热用户量大时Redis连接数和延迟都上来了。后来改成Flink的RedisSink批量Pipeline写入,单次请求携带多个字段,延迟降到几十毫秒。要注意TTL设计的坑:如果给每个用户特征都设了6小时TTL,大促期间活跃用户特征被删掉后,推荐服务会拿不到特征,必须配合推荐侧的降级策略,比如没有实时特征就退化为用类目Top商品或最近一次快照特征。
4. 自定义Data Source与Data Sink:连接器不够用时的正确姿势
4.1 什么时候需要自定义
不少刚开始接触Flink的同学会以为Flink自带所有连接器,但其实生产环境经常要自己写。我们遇到的需求有两类:一是公司内部自研的消息中间件或采集协议,官方Kafka连接器接不上;二是数据格式特殊,比如客户端上报的是二进制压缩格式,或者需要先解密再解析的业务流量。
自定义Source的本质是实现一个SourceFunction或RichParallelSourceFunction,关键点是并行度、Checkpoint和背压。如果你实现的Source没有正确保存消费位点,Flink任务重启后就会重复消费或丢数据,这在推荐场景里会直接导致特征计数错乱。我们的经验是:无论底层系统支持不支持,都要在Source里把"位点"存进ListState,从Checkpoint恢复时先读取位点,再决定从哪个位置继续消费。
4.2 一个自定义Source的完整片段
举个例子,有个合作方的事件数据通过自研的HTTP长轮询接口提供,不能直接用Kafka接入。我们写了一个HttpLongPollSource,继承RichParallelSourceFunction,启动时从状态里恢复上次的游标cursor,然后循环调用接口拉数据。拉回来的JSON先做字段校验,解析成统一的BehaviorEvent再发往下游。
这个Source里必须要处理两件事:第一是背压反馈,如果下游处理不过来,SourceFunction不能无限拉取,需要阻塞等待。我们通过collect方法天然实现了这一点,但要注意拉取批量不能太大,否则一次collect几千条会把算子内存打爆。第二是重试与退避,HTTP接口抖动是常态,没有重试策略的Source一遇到超时就整任务失败,我们设了指数退避重试,最多重试3次。
public class HttpLongPollSource extends RichParallelSourceFunction<BehaviorEvent> { private volatile boolean running = true; private transient ListState<String> cursorState; private String cursor; @Override public void run(SourceContext<BehaviorEvent> ctx) throws Exception { while (running) { List<BehaviorEvent> events = httpClient.poll(cursor); for (BehaviorEvent event : events) { ctx.collect(event); } cursor = updateCursor(cursor, events); cursorState.update(cursor); Thread.sleep(BACKOFF_MS); } } }4.3 一个自定义Sink的实现要点
自定义Sink最常见的坑是"写了但没生效""重复写""连接没释放"。我们写过一个到自研推荐指标平台的Sink,核心是批量攒数据、异步发送、失败重试。RichSinkFunction里用ListState缓存待发送数据,达到一定数量或时间后批量发送,发送成功才清理缓存。
幂等性必须提前设计。我们给每条特征数据分配一个唯一的事件ID,目标系统按事件ID去重。如果没有这个机制,任务重启后从Checkpoint恢复,Sink会重复发送最近一批数据,指标统计和特征覆盖都会出错。另外一个细节是连接复用,不要在invoke方法里每次新建连接,必须在open方法里初始化连接池或长连接,否则高并发下连接数会爆炸。
4.4 从词频统计理解Source和Sink的协作
很多初学者对Flink的第一印象就是"词频统计"。说实话,把Source、FlatMap、KeyBy、Sum、Sink这套链路跑通了,对理解分布式流处理帮助很大。我们内部带新人也是从"实时词频统计初体验"开始的。先通过自定义Source模拟从Socket、从文件、从随机数生成器读取数据,再用map和reduce统计窗口内的高频词,最后自定义Sink把结果输出到控制台或Redis。
这个项目的深层价值是,让新人理解Flink的数据流模型:数据不是存在某个全局表里,而是在算子之间流动;窗口不是定时器,而是根据Watermark和事件时间推进的;Checkpoint不是日志文件,而是让整个DAG具备"精确一次"恢复能力的机制。这些概念想通了,后面写复杂的推荐特征工程才不会犯方向性错误。比如窗口触发条件、状态清理策略、背压处理,本质上都能从词频统计这个小项目里找到原型。
5. 从Flink到TensorFlow Serving:在线推理链路落地的重点
5.1 模型导出:从训练到SavedModel的完整流程
模型训练好之后,导出成Serving能识别的格式是第一步。我们在导出时踩过最大的坑是"导出时没有把特征预处理逻辑打包进去"。比如模型在训练时对特征做了标准化,但标准化用的均值和方差存在训练代码里,导出模型时只导出了网络权重,没有把标准化算子加进去。结果是线上打分用的特征全是"裸特征",分数分布异常。
正确的做法是把预处理逻辑也写进tf.function里,让SavedModel的签名输入直接是原始特征,输出直接是打分结果。我们在导出时定义了一个serving_predict的签名,输入是item_id、user_behavior_seq、stats_features等原始字段,在函数内部完成标准化、缺失值填充、Embedding查表,最后输出score。这样Serving端拿到的是完整的模型,而不是需要客户端预处理的半成品。
@tf.function(input_signature=[tf.TensorSpec(shape=[None], dtype=tf.string)]) def serving_predict(serialized_examples): features = tf.io.parse_example(serialized_examples, feature_schema) # 预处理:标准化、填充缺失值 normalized = normalize(features) logits = model(normalized, training=False) return {"scores": logits} tf.saved_model.save(model, export_dir, signatures={"serving_predict": serving_predict.get_concrete_function()})5.2 gRPC与REST的选择:为什么线上用gRPC
TensorFlow Serving同时支持REST和gRPC两种接口。REST的好处是调试方便,用curl就能直接测。但线上推理我们用gRPC,原因很现实:第一,性能更好,protobuf二进制传输比JSON体积小很多,序列化开销低;第二,浮点数精度可控,JSON里的浮点数在序列化时可能被转成字符串,精度丢失,直接影响打分。第三,gRPC支持连接复用,长连接模式下请求耗时更稳定。
gRPC调Serving的流程是:先用protobuf构造tf.Example或PredictRequest,把特征填进去,然后通过PredictionServiceStub调用。Java端的实现要注意ManagedChannel的生命周期,不要每个请求都新建一个Channel,否则连接数和延迟都会出问题。我们是在推荐服务启动时初始化一个Channel,设置连接空闲超时和重连策略,这样压测时P99延迟能稳定在30ms左右。
syntax = "proto3"; message FeatureProto { map<string, float> float_features = 1; map<string, int64> int_features = 2; map<string, string> string_features = 3; }5.3 特征对齐:训练-推理一致性的核心保障
这一步是我们在项目里吃过最多亏的地方。最开始上线时模型离线AUC很高,线上点击率却上不来,查了很久发现是特征对齐问题。比如"用户历史点击商品的平均价格"这个特征,离线训练时用的是行为发生后一天的回溯统计,而在线推理时Flink算的是实时累积值,两者口径不一样。还有"近1小时行为数",离线脚本和Flink作业对"近1小时"的窗口边界定义不同,导致同一用户同一时刻算出来的特征值对不上。
解决方案是用一套"特征口径文档"约束所有实现,并且建立自动化校验任务。我们在Hive里定期跑一批用户的特征结果,同时在Redis里抽样相同用户的最新实时特征,做对比报表,一旦发现同一特征的平均绝对误差超过阈值,立刻报警。这个流程从上线到现在,帮我们抓到了好几个隐蔽的口径不一致问题。另一个建议是,把特征版本号写进Redis的key里,比如user_feature_v3_12345,这样模型和特征可以各自独立版本化,回滚和压测都方便。
5.4 降级方案:模型超时、特征缺失时推荐什么
降级不是一个"要不要做"的问题,而是"怎么做才优雅"的问题。我们设计了三层降级。第一层是Serving超时或返回异常,推荐服务退化为规则策略,按商品热度、用户最近浏览类目Top商品、运营配置的固定位商品来填充。第二层是用户实时特征不存在,比如新用户或者Redis过期,这时用离线画像特征和商品热门度打分。第三层是整个Serving集群不可用,直接返回运营预置的推荐结果,保证页面不空白。
降级里有一个容易被忽略的细节:一定要在日志和监控里标记降级原因。我们给每次推荐响应都加了degrade_reason字段,在监控大盘上按原因聚合,能清楚看到是模型超时还是特征缺失导致的降级。如果特征缺失的比例过高,说明Flink链路或Redis TTL配置有问题,要优先修;如果模型超时比例高,说明Serving的资源配置不足,要扩容。
6. 线上真实踩坑:JDBC连接器异常、Hive Sink入不了表、CDC部署和火焰图
6.1 Flink JDBC连接器为什么总爆"connection closed"
这个异常我们在维表关联阶段遇到过很多次,表现形式是SQLException: Communications link failure或者Connection is closed。刚开始以为是MySQL端把连接断了,后来抓包排查才发现,根子是Flink的JDBC连接器内部用的连接池空闲连接超时被MySQL服务端回收了,但连接池没感知,下次请求还用这个坏连接,自然报错。
排查链路可以复现一下。第一步,看Flink TaskManager日志里报错的时间点,会发现集中在低流量时段,比如凌晨或大促间隙。第二步,查MySQL端的wait_timeout设置,默认一般是8小时,但Flink连接池的maxLifetime如果设置超过这个值,空闲连接就会被服务端杀。第三步,确认连接器版本,旧版本的JDBC连接器对连接有效性检查做得不够,直接复用坏连接。
解决的组合拳是:把连接池的maxLifetime调到小于MySQL的wait_timeout,加上连接有效性检测SQL(validationQuery=SELECT 1),并开启自动重连。另外,维表join的查询量如果很大,不要在JDBC连接器里做复杂SQL,尽量简化成按主键查询,把复杂join拆到Flink里做。
6.2 Sink到Hive表数据不入表的问题
我们的特征Topic要落到Hive表供离线训练使用,遇到过任务运行正常但查询Hive表一直没数据的情况。首先是确认Flink作业的Sink是否真的在提交数据。打开TaskManager日志,看到open了Hive连接,但一直没有触发commit动作,说明数据还在buffer里,没有达到提交阈值。
根因是Flink StreamingFileSink或HiveSink的提交策略默认是"分批提交",数据要攒够一定大小或经过一定时间间隔才commit。测试环境数据量小,一直没触发提交条件,看起来就像没数据。解决办法是把auto-compaction打开,或者调低提交间隔参数,比如sink.partition-commit.policy.kind=success-file、sink.partition-commit.trigger=partition-time,让数据及时可见。
还有一个隐藏很深的坑是分区字段的值。Hive分区字段如果用的时间函数,和Flink的Watermark时区不一致,会导致数据被写进了错误的分区,看起来"没数据"其实是数据去了别的时间分区。我们最后统一用Flink的时间函数+固定UTC偏移来生成分区路径,再也没有出现"不入表"的诡异问题。
6.3 Flink CDC Pipeline部署与OpenMetadata血缘管理
导购平台另一个重要链路是数据库变更的实时捕获。我们用Flink CDC把MySQL业务库的变更同步到数仓,支撑商品信息更新、库存变更等场景。CDC Pipeline部署比普通Flink作业多几个注意点:第一是连接器要用com.ververica.cdc.connectors.mysql.source.MySqlSource,并且要配置startupOptions,否则任务重启后会从binlog最末尾开始读,中间的数据全丢了。第二是并行度不能随意调整,CDC Source是单线程读binlog的,如果并行度大于1,要确认上游MySQL的binlog格式和表结构是否支持多并行度。
数据血缘这块我们用了OpenMetadata来管理Flink作业的元数据。好处是能快速回答"这个特征字段是从哪张表来的""这个Hive分区被哪些Flink作业写入"这类问题。但要注意,OpenMetadata对Flink SQL的血缘解析比较依赖Calcite优化器,自定义的DataStream算子它看不到,需要在OpenMetadata里手动维护作业级元数据,把Source、Sink的上下游信息补全。我们实际使用中,血缘解析的准确率大概在八成左右,剩下的要靠人工规则补充。
6.4 火焰图排查吞吐瓶颈
实时特征作业上线后,我们发现高峰期吞吐上不去,但CPU使用率没有跑满。用JFR和Async Profiler抓了火焰图,很快定位到热点:JSON解析占用了接近40%的CPU。我们用的是Jackson的JsonNode树模型去解析埋点数据,每个字段都要创建Node对象,在每秒钟几万条消息的规模下GC压力特别大。后来改成直接基于Jackson的流式API解析,并尽量复用ObjectMapper,CPU占用直接降了一半。
另一个火焰图暴露的问题是正则表达式。清洗URL和商品ID时用了几条复杂正则,在高频匹配时非常消耗CPU。改造成简单的split和字符串判断后,吞吐提升了30%。火焰图的价值就在这里,不用猜,直接看哪个方法栈在最上面。我们养成了习惯:每次调优之前先抓火焰图,改动后再抓一张对比,用数据说话。
7. 性能压测、延迟监控与上线后的持续治理
7.1 端到端延迟的监控口径与实现
端到端延迟是这套架构里最重要的观测指标。我们埋点的实现方式是:用户行为发生时就记录event_time,这条数据进入Kafka和Flink处理时都保留这个字段,特征写入Redis时同时写入feature_time。推荐服务读取特征时,用当前时间减去feature_time,得到"特征延迟"。监控大盘上按秒级分组展示延迟分布,一旦P95超过5秒,就触发告警。
这里有个很关键的细节:Flink处理本身产生的延迟很小,真正的延迟大头在埋点上报和Kafka消费积压。所以要分开监控三个指标:上报延迟、消费延迟、处理延迟。排查问题的时候,先看Kafka的ConsumerLag,再往上游查埋点SDK,最后才看Flink作业内部的算子耗时。我们曾经把一个"推荐位更新慢"的问题定位到客户端SDK的批量上报策略上,而不是Flink,就是因为分开监控帮了大忙。
7.2 压测方案与容量预估
压测是上线前绕不开的一步。我们的方案分两层:第一层是单链路压测,分别压Flink作业和TensorFlow Serving;第二层是端到端压测,用压测工具模拟真实用户行为流量,打到推荐服务,观察全链路表现。单链路压测时,Flink作业的吞吐我们用5000事件每秒、10000事件每秒、20000事件每秒三档递进,观察算子处理耗时和背压指标。TensorFlow Serving的压测则用gRPC客户端并发打请求,看P99延迟和GPU利用率。
容量预估的经验公式是:推荐高峰QPS、行为事件放大倍数、特征更新频率三者相乘,得到特征写入Redis的TPS,再乘以一定冗余系数。比如线上高峰期推荐QPS是8000,RT是30ms,那么同时刻产生的行为事件大约是每秒12000条,Flink侧需要处理的特征更新约为每秒10000次。这样算出来的并行度基本够用,再多留30%的余量应对大促。
7.3 收益验证与持续优化
推荐系统的最终评价必须回到业务指标。我们上线实时特征后做了两组AB实验:一组是"离线特征+离线模型"作为基线,一组是"实时特征+在线推理"作为实验组。实验组首页推荐点击率提升了18%左右,加购转化率提升9%左右。但这里也要提醒,不同导购场景收益差别很大,搜索结果页因为用户意图已经明确,实时特征的增益主要在"意图遗忘"上;首页和详情页才是实时特征的主战场。
持续优化上,我们每两周做一次特征增量迭代,最有效的两个方向是:扩展行为序列长度(从50扩展到200,点击率又涨了接近5%)、加入实时上下文特征(如当前是否在大促会场、商品实时库存状态)。模型侧也在试多任务学习,把点击、加购、下单几个目标一起训练,后续计划把在线学习接进来,让模型参数也能随着实时数据流更新。这一步比特征工程复杂很多,还在实验阶段。
写在最后:这套架构里,比技术更重要的两件事
做完整套架构,我自己的体会是,技术选型其实不是最难的部分,Flink、TensorFlow Serving都是非常成熟的开源组件,重点是团队的"特征口径一致性"意识和"降级预案完备性"。实时特征再快,如果和离线训练的口径对不上,模型打分就是乱来;在线推理再高效,如果降级策略设计不好,一次抖动就能让用户看到空白推荐位。
还有一点经验想分享给正在做类似项目的朋友:不要一开始就把架构搞得太重。我们早期其实只用Flink做了最核心的"近30分钟行为聚合",模型还是离线的,先跑通再逐步加行为序列、加在线推理、加CDC。每一步上线都做AB对比,确认有效再继续。这样不仅风险小,团队对每一块组件的理解和掌控度也会扎实很多。推荐系统没有银弹,实时化只是把"数据到决策"的距离缩短了,真正的价值还是来自对业务场景的理解和对细节的较真。