Flink实时特征工程与TensorFlow Serving在线推理实战:推荐系统秒级响应架构复盘

发布时间:2026/9/30 22:07:00
Flink实时特征工程与TensorFlow Serving在线推理实战:推荐系统秒级响应架构复盘
最近在帮朋友优化一个导购平台的推荐链路感触很深。很多团队做推荐模型训练和特征工程都挺完善但到了线上用户从点击商品到推荐位更新之间隔着十几分钟甚至更久转化率就被拖垮了。我们后来落地了一套基于Flink的实时特征工程 TensorFlow Serving在线推理架构把推荐引擎的时效性从分钟级拉到了秒级。这篇内容不是PPT架构图是我们在真实业务里踩坑、调优、上线后的完整复盘包括自定义Source和Sink的实战写法、JDBC连接器异常、CDC部署、Hive Sink入不了表这些具体问题希望对做推荐和实时计算的同行有参考价值。当时我们明确的目标很朴素用户刚搜过冬季连衣裙回到首页推荐位必须能看到相关商品用户加购了某品牌吹风机详情页的相似好物要尽快跟着变。要做到这一点光靠离线T1的特征根本无法完成必须在事件发生的瞬间把行为特征算出来再扔给模型打分。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用ValueStateListBehaviorItem保存最近N个行为用MapStateString, Long记录各行为类型的计数。每来一条新行为先判断当前时间与上一条行为的间隔如果超过30分钟说明是新会话要重置会话内状态。然后把行为压入列表如果超过200条就剔除最旧的并更新时间衰减因子。最后把编码后的特征向量发往下游。这里有两个特别重要的细节。第一状态大小必须收敛。如果无限保留用户所有历史行为大促期间热用户的状态会暴涨直接拖垮内存和Checkpoint。我们的做法是ListState只保留最近200条每处理完一批就清理超出的部分同时配合StateTtlConfig给状态设置24小时的TTL。第二更新下游特征时要有版本号。Redis里同一个用户key可能被并发写入需要记录一个事件时间戳下游读取特征时能感知到数据新鲜度而不是盲信任一个旧快照。// 伪代码展示核心状态管理逻辑 public class UserBehaviorSeqProcessFunction extends KeyedProcessFunctionString, BehaviorEvent, UserFeature { private transient ListStateBehaviorItem recentBehaviors; private transient ValueStateLong lastEventTime; Override public void processElement(BehaviorEvent event, Context ctx, CollectorUserFeature 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 维表缓存实时特征不只有行为序列还要关联上商品信息、用户画像、价格区间等维度属性。如果每来一条行为都同步查一次MySQLFlink的吞吐量会直接被打死。我们用了Async I/O 本地缓存两层方案。Async I/O是Flink官方提供的异步访问外部系统的API核心优势是发一个异步请求后不阻塞算子线程等回调返回后再继续处理。我们在AsyncFunction里并发查询用户画像和商品信息查询结果拼接到特征上。但异步也只是解决了并发度问题没解决热点key的重复查询问题。所以我们又加了一层CacheKey, 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 RichParallelSourceFunctionBehaviorEvent { private volatile boolean running true; private transient ListStateString cursorState; private String cursor; Override public void run(SourceContextBehaviorEvent ctx) throws Exception { while (running) { ListBehaviorEvent 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], dtypetf.string)]) def serving_predict(serialized_examples): features tf.io.parse_example(serialized_examples, feature_schema) # 预处理标准化、填充缺失值 normalized normalize(features) logits model(normalized, trainingFalse) return {scores: logits} tf.saved_model.save(model, export_dir, signatures{serving_predict: serving_predict.get_concrete_function()})5.2 gRPC与REST的选择为什么线上用gRPCTensorFlow 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 { mapstring, float float_features 1; mapstring, int64 int_features 2; mapstring, 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加上连接有效性检测SQLvalidationQuerySELECT 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.kindsuccess-file、sink.partition-commit.triggerpartition-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解析并尽量复用ObjectMapperCPU占用直接降了一半。另一个火焰图暴露的问题是正则表达式。清洗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是8000RT是30ms那么同时刻产生的行为事件大约是每秒12000条Flink侧需要处理的特征更新约为每秒10000次。这样算出来的并行度基本够用再多留30%的余量应对大促。7.3 收益验证与持续优化推荐系统的最终评价必须回到业务指标。我们上线实时特征后做了两组AB实验一组是离线特征离线模型作为基线一组是实时特征在线推理作为实验组。实验组首页推荐点击率提升了18%左右加购转化率提升9%左右。但这里也要提醒不同导购场景收益差别很大搜索结果页因为用户意图已经明确实时特征的增益主要在意图遗忘上首页和详情页才是实时特征的主战场。持续优化上我们每两周做一次特征增量迭代最有效的两个方向是扩展行为序列长度从50扩展到200点击率又涨了接近5%、加入实时上下文特征如当前是否在大促会场、商品实时库存状态。模型侧也在试多任务学习把点击、加购、下单几个目标一起训练后续计划把在线学习接进来让模型参数也能随着实时数据流更新。这一步比特征工程复杂很多还在实验阶段。写在最后这套架构里比技术更重要的两件事做完整套架构我自己的体会是技术选型其实不是最难的部分Flink、TensorFlow Serving都是非常成熟的开源组件重点是团队的特征口径一致性意识和降级预案完备性。实时特征再快如果和离线训练的口径对不上模型打分就是乱来在线推理再高效如果降级策略设计不好一次抖动就能让用户看到空白推荐位。还有一点经验想分享给正在做类似项目的朋友不要一开始就把架构搞得太重。我们早期其实只用Flink做了最核心的近30分钟行为聚合模型还是离线的先跑通再逐步加行为序列、加在线推理、加CDC。每一步上线都做AB对比确认有效再继续。这样不仅风险小团队对每一块组件的理解和掌控度也会扎实很多。推荐系统没有银弹实时化只是把数据到决策的距离缩短了真正的价值还是来自对业务场景的理解和对细节的较真。