简介:基于Flink的商品实时推荐系统项目包,适合大数据开发与推荐系统学习者,用于解决实时商品热度统计、用户画像构建及个性化推荐排序等核心问题。资源共109个文件、压缩包3.74MB,以68个Java源文件为主,另有SQL建表脚本、HBase建表语句、XML配置、Properties与YML配置文件、Kafka模拟数据脚本及HTML展示页面,覆盖从数据接入、缓存存储到推荐计算和前端展示的完整链路。系统采用Flink统计商品热度并写入Redis缓存,同时分析用户日志生成画像标签、实时记录存入HBase;用户发起推荐请求时,系统基于用户画像对热度榜重排序,并融合协同过滤与标签推荐两个模块为每个商品补充关联产品,最终返回个性化推荐列表。目前已有130人学习。通过源码可掌握Flink流处理、Redis与HBase集成、协同过滤及标签推荐算法的工程落地方法,适合进阶实战参考。
1. 实时推荐不是算法竞赛,而是特征时效性的竞赛
用户在商品详情页刚滑走的第 5 秒,离线推荐系统还在等当天 T+1 日志清洗完成;而基于 Flink 的商品实时推荐系统,已经把“刚刚浏览过这件商品”写成了下一轮召回的特征。这套方案要解决的核心问题,是把用户行为从小时级延迟压缩到秒级,让每一次点击都能真实影响下一次推荐内容。它不引入什么新模型,而是把从埋点到特征、再到召回排序的整条链路做成实时闭环。适合已经有埋点、数据仓库和离线推荐,但特征延迟仍停留在小时级以上的团队;新手可以把这里当作 flink 实时计算的第一个完整项目,熟手则能看到状态管理和连接器选型的真实边界。
2. 整体架构与选型:Flink 在推荐链路里的准确位置
实时推荐不是把离线推荐的代码换个执行引擎就跑,而是要把原来“按天算好写回表”的特征,改成“持续流入、持续更新”的流式特征。Flink 在这里承担的是特征计算层,而不是完整的推荐系统。先把这层位置摆正,后面每一步选型才有依据。
2.1 推荐系统的四层链路:埋点、特征、召回、排序
推荐系统看起来是一个在线服务,实际是一条多层接力链路。第一层是埋点采集,App 或 Web 页面把用户行为(曝光、点击、加购、下单)实时上报到 Kafka;第二层是流式特征计算,这是 Flink 的主战场,按用户维度聚合行为序列,按物品维度计算关联关系;第三层是召回,在线服务根据用户当前行为从特征库里取候选集;第四层是排序,把候选集按预估点击率排序后输出到页面。
中间第二层之所以必须用流式计算而不是离线批处理,是因为召回和排序都要读“上一秒”的特征。离线推荐把特征打成 T+1 宽表,在线服务读到的永远是一天前的偏好;实时推荐把这个宽表拆成不断追加的流式特征,用户只要在窗口期内发生了行为,特征就被更新,下游立刻可见。商品信息这类变化不频繁的维度数据,常见做法是用 Flink CDC Pipeline 把业务库的变更直接同步到特征存储,避免定时全量同步带来的延迟黑洞。需要提醒的是,flink cdc 安装部署本身不难,但 pipeline 模式对连接器版本兼容性比较敏感,数据量大之前先花半天在测试环境跑通整条链路,比上了生产再排查划算得多。
2.2 为什么选 Flink 而不是 Spark Streaming 或自研管道
落到选型,团队里最常见的争论是 Spark Streaming 也能消费 Kafka,为什么非要 Flink。在实时推荐这个场景里,差距主要体现在三个点:事件时间语义、状态管理和容错恢复。埋点日志在 Kafka 里只要稍微积压,基于处理时间的窗口统计就会把 10 点的事件算进 10 点 05 分的窗口,推荐结果看起来在更新,实际上是错位的。Flink 的事件时间加水印机制可以按日志里的时间戳重新对齐,行为序列不会因为中间链路抖动而错乱。
| 对比维度 | Spark Streaming | Flink |
|---|---|---|
| 时间语义 | 以处理时间为主,事件时间支持有限 | 事件时间 + 水印原生支持 |
| 状态管理 | 窗口内状态,窗口结束即释放 | ValueState / ListState + TTL |
| 容错恢复 | 批式重跑,存在重复计算 | 分布式快照,可做到精确一次 |
| 端到端延迟 | 秒级到分钟级 | 秒级以内 |
这张对比表基本就是拍板依据。实时推荐属于典型的“状态密集 + 窗口密集”场景,每个用户的行为序列是持续累积的状态,每个窗口的共现统计依赖精确的时间对齐,这两点正好都是 Flink 的长项。部署层面,flink 安装配置到部署按标准集群流程走就可以,但要记住:环境搭建只决定作业能不能跑起来,窗口和状态的设计才决定推荐效果好坏。
2.3 数据流与角色划分:从 Kafka 到 Redis 的接力
在推荐场景里,Flink 作业通常拆成两个角色。一个是用户特征作业,消费行为日志,按 userId 分组,用状态维护最近 N 分钟的行为序列,产出的结果写入 Redis;另一个是物品关联作业,消费同一份行为日志,计算物品之间的共现关系,产出物品相似度集合。两个作业可以合并成一个 DAG,也可以分开部署,我一般建议分开。二者的吞吐特征和故障影响面不一样,拆开后可以独立调优并行度、独立重启,排查问题时也不需要把整条链路停掉。
这里要特别说一句:Flink 自带的 Kafka Source 和 Redis Sink 往往不够用。公司自研埋点 SDK 的上报格式、内部 Redis 集群的鉴权方式,都需要自定义 data source 与 data sink,这也是 flink 实时计算进阶篇里最常被问到的部分。自定义 Source 的核心是理清消费位点、反压信号和 schema 解析三件事;自定义 Sink 的核心则是批量写入、失败重试和幂等。如果公司已经上了 openmetadata 这类数据血缘平台,Flink 作业的血缘关系可以被自动采集,字段级溯源到 Kafka topic,这对后面排查“特征被谁改过”很有价值,建议在一开始就把作业名和 topic 命名规范定好,血缘平台才能真正发挥作用。
3. 从埋点到特征:一套最小可复现的实时计算作业
看完架构,下一步是把用户特征作业落地。下面这套逻辑是实时推荐链路中最基本、也是最先要跑通的部分。它做的事情可以概括为:接收行为日志,按用户维护最近一段时间的行为序列,把这个序列作为后续召回和排序的输入。
3.1 行为数据的字段规范:决定后续所有计算的上限
行为日志字段规范决定特征质量的上限,Flink 作业只是把规范变成结果。一条点击事件至少要有这些字段:userId、itemId、behaviorType(click、cart、order)、sceneId(首页、详情页、搜索页)、ts(事件发生时间,毫秒时间戳)。如果公司埋点已经存在,最好在接入层补齐,不要指望 Flink 侧做数据清洗,字段缺失只能靠默认值硬填,后续聚合和 join 会全部偏离。
一个常见误区是,为了兼容所有业务线,把行为日志设计成几十个字段的大宽表。实时推荐真正高频使用的字段通常不超过八个,字段越多,Kafka 序列化开销和 Flink 反序列化开销越大。我一般处理方式是 topic 里只保留推荐链路必用字段,其余字段走旁路给数据仓库,这样 Kafka 吞吐和 Flink 作业性能都能保住。
3.2 用 KeyedProcessFunction 维护用户实时行为序列
下面的代码是用户特征作业的核心骨架。它消费 Kafka 行为日志,按 userId 分组,用 ListState 保存用户最近 30 分钟的行为序列,并用事件时间定时器在 30 分钟后清理状态。
DataStream<UserBehavior> stream = env .addSource(new FlinkKafkaConsumer<>("user_behavior", new JSONDeserializationSchema(), kafkaProps)) .assignTimestampsAndWatermarks( WatermarkStrategy.<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, timestamp) -> event.getTs())) .keyBy(UserBehavior::getUserId) .process(new KeyedProcessFunction<String, UserBehavior, UserRecentItems>() { private transient ListState<UserBehavior> recentItems; private transient ValueState<Long> cleanupTimer; @Override public void open(Configuration parameters) { ListStateDescriptor<UserBehavior> descriptor = new ListStateDescriptor<>("recentItems", UserBehavior.class); recentItems = getRuntimeContext().getListState(descriptor); cleanupTimer = getRuntimeContext().getState(new ValueStateDescriptor<>("cleanupTimer", Long.class)); } @Override public void processElement(UserBehavior value, Context ctx, Collector<UserRecentItems> out) throws Exception { recentItems.add(value); if (cleanupTimer.value() == null) { long timerTime = ctx.timestamp() + 30 * 60 * 1000L; ctx.timerService().registerEventTimeTimer(timerTime); cleanupTimer.update(timerTime); } out.collect(buildOutput(recentItems.get())); } @Override public void onTimer(long timestamp, OnTimerContext ctx, Collector<UserRecentItems> out) throws Exception { recentItems.clear(); cleanupTimer.clear(); } });逻辑说明:这段代码的关键不是 map 和 filter,而是状态和时间语义。assignTimestampsAndWatermarks 告诉 Flink 使用事件时间而非处理时间,后续窗口和定时器都基于日志里的 ts 推进;KeyedProcessFunction 里每个用户维护一份独立的 ListState,保存的是行为对象列表而不是拼好的字符串,因为后续召回阶段还要从这些对象里取 itemId 和 sceneId;定时器注册在第一条事件的事件时间上加 30 分钟,时间到达后触发 onTimer 清理状态,确保长期不活跃用户的状态不会积压。
参数说明:forBoundedOutOfOrderness 的 10 秒是乱序容忍度,要根据埋点上报链路的实际延迟调整。设得太大,窗口出结果变慢,特征时效性下降;设得太小,迟到数据被丢弃,特征不完整。30 分钟的行为序列窗口决定“实时偏好”的敏感度:做闪购类业务我一般压缩到 5 分钟,做常规电商可以放宽到 30 到 60 分钟。ListState 可以额外配置 TTL,按天或按小时清理,防止序列无限增长。还需要注意,如果只写状态而不开 checkpoint,进程重启后状态全部丢失,实时特征就断了。checkpoint 是实时作业的后悔药,建议一上生产就打开。
3.3 ItemCF 的实时化:物品相似度不能全量算
离线 ItemCF 的做法是统计用户行为序列中物品对共现次数,再计算余弦相似度。实时化之后,最朴素的思路是把全量相似度计算改成“增量共现窗口”,只在用户最近点击列表里生成物品对。下面这段代码是核心示意,Flink 版本不同 API 命名会有些差异,但整体结构是通用的。
// 过滤出点击行为,按用户分组维护最近 N 个点击物品 DataStream<UserBehavior> clicks = stream.filter(e -> "click".equals(e.getBehaviorType())); clicks.keyBy(UserBehavior::getUserId) .process(new KeyedProcessFunction<String, UserBehavior, ItemPair>() { private transient ListState<String> itemSeq; @Override public void open(Configuration parameters) { itemSeq = getRuntimeContext() .getListState(new ListStateDescriptor<>("itemSeq", String.class)); } @Override public void processElement(UserBehavior value, Context ctx, Collector<ItemPair> out) throws Exception { itemSeq.add(value.getItemId()); List<String> items = new ArrayList<>(); for (String item : itemSeq.get()) { items.add(item); } // 只保留最近 20 个点击,避免列表无限膨胀 if (items.size() > 20) { items = items.subList(items.size() - 20, items.size()); itemSeq.clear(); itemSeq.addAll(items); } // 把最后一个点击和它之前的所有点击组成共现对 if (items.size() >= 2) { String last = items.get(items.size() - 1); for (String prev : items.subList(0, items.size() - 1)) { out.collect(new ItemPair(prev, last)); } } } });逻辑说明:这个算子把每个用户的点击行为按时间顺序累积,每次新点击到达时,都和当前序列里已有的物品组成共现对。比如用户依次点击了 A、B、C,那么在 C 到达时会输出 C-A、C-B 两对。这样做避免了离线方式下全量扫描用户历史的开销,只关注“最近发生了关联”的物品。
参数说明:保留最近 20 个点击,是召回质量和计算量的折中。20 太小,长尾关联丢失;20 太大,每个用户每次点击都要生成更多物品对,下游共现统计压力成倍增长。生产环境里还可以再加一层处理:生成的 ItemPair 进入滑动窗口按对计数,比如 30 分钟窗口每 5 分钟滑动一次,只保留出现次数超过阈值的物品对写入 Redis,这样能过滤掉误触带来的噪声共现。热门物品的相似集合要注意长度控制,我一般每个 itemId 最多保留 50 个相似品,避免头部商品把长尾商品的曝光全部吃掉。
3.4 特征结果落地:Redis 和 HBase 的选型与参数
特征计算完成之后,在线服务要能高效读取。用户行为序列这类“每个用户一条”的数据,用 Redis 最合适;物品相似度矩阵这类“每个物品多条”的数据,也适合放 Redis。QPS 高、数据量在百万级以内,Redis 完全扛得住;如果用户量上亿、序列特征总量超过 Redis 内存的合理范围,就要把序列特征落到 HBase,行键设计成 userId 反转加时间戳前缀,避免热点分片。
写 Redis 时要注意 key 的过期策略。行为序列和物品相似集合都不需要永久保存,给 key 设置 TTL 可以防止存储无限膨胀。
| 数据 | key 模式 | 结构 | 建议 TTL |
|---|---|---|---|
| 用户最近行为 | user:recent:{userId} | ZSet 按时间戳排序 | 24 小时 |
| 物品相似集合 | item:sim:{itemId} | Hash 存相似度得分 | 6 小时 |
| 实时热销榜 | rank:hot:overall | ZSet 按曝光/点击加权 | 10 分钟 |
并行度和 Redis 分片也需要对齐。如果写入用的并行子任务数大于 Redis 集群分片数,单个分片会过热;我一般会在写入前对存储 key 做一次 keyBy 预聚合,再直连 Redis。项目初期如果只是照着 flink 菜鸟教程里的词频统计做一遍,对理解算子有帮助,但真正跑推荐作业时,状态、水印、窗口三个概念才是主战场,词频统计里 map 加 sink 的三件套支撑不起这套链路的复杂度。
4. 实时推荐避坑记录:五个翻车现场的根因与解决办法
实时推荐链路长,每一层都可能出问题。这里整理的是最常见的五个翻车场景,每一条都按现象、原因、解决三步写清楚,照着排查可以省下大量抓日志的时间。
4.1 事件时间与处理时间混用,窗口数据整体漂移
现象:Kafka 里积压了 10 分钟日志后,窗口统计出的“最近 5 分钟热门商品”看起来在更新,但和页面实际行为对不上,热门商品整体滞后。
原因:作业里有的窗口用处理时间,有的用事件时间。处理时间以 Flink 机器本地时钟为准,日志在 Kafka 积压后,处理时间比事件发生时间晚,导致两个窗口的时间基准不一致,数据自然漂移。
解决:全作业统一使用事件时间,Kafka Source 配置 WatermarkStrategy,上游 topic 的延迟会直接体现在 watermark 里。运维上要关注消费组 Lag 指标,Lag 持续上涨说明消费能力不足,先解决这个再谈特征准确。排查时可以用 flink 火焰图定位是哪个算子耗时异常,先确认瓶颈在 source、窗口计算还是 sink,再决定加并行度还是优化序列化。
4.2 维表关联报错:flink 的 jdbc 连接器异常
现象:Flink SQL 作业里 join 商品维表,运行几个小时后开始报 flink 的 jdbc 连接器异常,连接被关闭,作业重启后恢复正常,过几个小时又复发。
原因:维表 join 的默认实现是每条数据都查一次数据库,Kafka 高峰流量下连接数被打满,超时后被连接池回收,后续请求拿到失效连接就开始报错。更隐蔽的是,连接池回收和 Flink 算子线程之间不同步,偶发异常容易被误判为网络问题。
解决:常见做法是用 lookup join 加缓存。Flink SQL 维表 DDL 里设置 lookup.cache.max-rows 和 lookup.cache.ttl,商品维表这类变化不频繁的维度,缓存 5 到 10 分钟即可;如果流量再大,就用 asyncio 改造维表查询,或者把维表预加载成广播状态。血泪经验是:维表 join 永远不要用同步逐条查询,无论单次查询多快,都扛不住 Kafka 高峰吞吐。
4.3 状态 TTL 设置不当,checkpoint 连环失败
现象:作业运行一段时间后,checkpoint 持续超时,连续失败触发自动重启,用户特征大面积丢失。
原因:ListState 里存的是完整行为对象,每个对象包含多个字段,TTL 设置成一天,数据量上来后状态体积膨胀,checkpoint 序列化的数据量超过默认超时阈值。
解决:状态量大的作业先评估存储格式,优先用紧凑的字节数组而不是完整 JSON 对象,字段能省则省;可以适当调大 checkpoint 超时时间,同时开启增量 checkpoint。状态体积调参有点像玄学,本质是把“保留多少历史行为”和“多久能备份完”两件事对齐,不是把 Flink 内存参数盲目调大。
4.4 Sink 到 Hive 表数据不入表,结果静默丢失
现象:Flink 作业显示正常结束,日志里也写了写入成功,但 Hive 表查不到数据,或者只有最后一批数据。
原因:常见于流式结果积累为小文件分批写入 Hive,并行度大于分区数时,多个并行子任务各自提交文件,元数据没有合并;另外流式写 Hive 如果没开分区提交,数据会一直停在 staging 目录,flink sink hive表数据不入表就是这么来的。
解决:Flink SQL 写 Hive 要显式配置分区提交参数,指定提交触发时机和文件格式;同时控制写 Hive 的并行度不要超过分区数,避免小文件碎片。排查时先看 HDFS staging 目录有没有数据:有数据说明是提交策略问题,没数据说明上游就没输出,这时候用火焰图看算子是 idle 还是 busy,再决定往上排查 Kafka 消费还是往下排查 sink。最忌讳的是只把 Hive 连接参数改一遍又重跑,浪费一个完整窗口周期。
4.5 迟到数据被静默丢弃,夜间特征缺失
现象:白天实时推荐表现正常,凌晨网络抖动导致部分行为日志延迟到达,第二天用户看到的推荐里丢失了昨晚的浏览偏好。
原因:水印设置了固定的 10 秒乱序容忍度,夜里上报链路抖动超过 10 秒后,迟到数据被窗口直接丢弃,用户特征里缺失了这部分行为。
解决:给窗口计算增加迟到侧输出,用 OutputTag 收集迟到的行为单独处理;同时把水印容忍度按业务时段调整。推荐系统对“最近一次行为”的依赖很强,宁可在极端场景下多算一次重复统计,也不能把迟到数据静默丢进黑匣子里。
5. 召回与排序:把特征变成用户看到的商品
Flink 算出的特征不会直接变成推荐结果,中间还要经过召回和排序。这一层的设计决定了实时特征到底能不能转化为页面上的点击率提升。
5.1 实时召回候选集的组装:一次读取、聚合去重
召回的核心是用“用户刚刚的行为”去特征库里取关联物品。在线服务收到推荐请求后,先读 Redis 里这个用户的最近行为序列,再批量查每个行为物品的相似集合,最后合并去重。这里最大的性能隐患是网络往返次数,不能用循环逐条请求。
// 在线召回服务核心逻辑:读取行为序列与相似物品集合 String recentKey = "user:recent:" + userId; Set<String> recentItems = jedis.zrevrange(recentKey, 0, 9); List<String> candidates = new ArrayList<>(); // 批量获取相似集合,避免逐条循环请求 for (String itemId : recentItems) { Map<String, String> similar = jedis.hgetAll("item:sim:" + itemId); candidates.addAll(similar.keySet()); } // 过滤已购商品,按相似度得分排序 candidates.removeAll(purchasedItems);逻辑说明:这段代码先在 Redis 里取用户最近 10 个行为物品,再逐个取相似集合。生产环境里要把 for 循环改成 pipeline 或 mget,减少网络 RTT;hgetAll 返回的是 Hash 结构的 field-value,field 是相似物品 ID,value 是相似度得分。排序时以相似度得分为主排序键,行为时间远近为副排序键,保证“最近看过的物品的相似品”排在更前面。
参数说明:zrevrange 取 10 个行为物品,是召回覆盖面和查询开销之间的折中;每个物品保留 50 个相似品,10 个物品就是最多 500 个候选。候选集过大时排序阶段的压力会明显变高,TP99 延迟上涨,需要配合实际压测来调。这里给出一个常用的初始参数表,上线时根据自己的流量调整。
| 参数 | 初始值 | 说明 |
|---|---|---|
| 召回宽度 | 10 | 用户行为序列取数量 |
| 相似集合长度 | 50 | 每个物品保留相似品上限 |
| 候选集上限 | 500 | 召回去重后的候选总数 |
| 排序输出条数 | 20 | 最终展示的商品数 |
5.2 轻量排序:行为加权和时间衰减的计算方式
召回给出候选商品后,排序决定最终展示的 20 个。实时链路里排序模型不宜太重,常见做法是先做规则加权,再叠加一个轻量逻辑回归。规则加权的核心是把行为类型价值差异和时间衰减写清楚:下单权重大于加购,加购权重大于点击,时间越近权重越大。常见的时间衰减公式是 weight 乘以 exp(-lambda * elapsedMinutes),lambda 取值决定衰减速度,做秒级时效类业务时 lambda 要调大。
排序特征的来源就是 Flink 已经算好的那套结果:用户最近 30 分钟的行为序列长度、候选物品的相似度得分、候选物品价格带与用户历史消费价格带的匹配程度、候选物品是否和用户最近浏览品类一致。这些特征拼成向量喂给逻辑回归,输出预估点击率。逻辑回归的训练样本可以用曝光日志和点击日志离线构造,不需要在线学习,训练链路和实时特征链路完全解耦,样本延迟一天也不会影响线上推理。
5.3 兜底策略:冷启动与热销榜的边界
新用户没有任何行为序列,召回阶段读 Redis 拿不到数据,候选集为空。这种情况直接返回热销榜。热销榜可以由同一个 Flink 作业在窗口内统计曝光和点击加权计算,写入 Redis 时用独立 key 和更短 TTL,避免和个性化特征混在一起。兜底策略还要考虑不同场景的差异:未登录用户的兜底应该是全局热销,登录但无行为的新用户可以叠加地域和品类偏好,搜索后无结果则需要回到搜索前场景的兜底列表。冷启动用户占比较高时,实时推荐整体的点击率指标会被明显拉低,这是正常现象,需要按用户状态分层看指标,不要只看全局平均值。
6. 上线前怎么验证这套推荐链路真的变好了
实时推荐链路很长,上线最怕的是直接切流量后指标下跌,又找不到原因。验证方式不需要一开始就上实验平台,旁路日志法足够解决问题。
6.1 旁路日志做离线 AB,不切线上流量也能评估
在线服务把“实时推荐给出的候选集、排序结果和最终曝光商品”落一份日志,再用离线脚本模拟“如果用户当时看到这个结果,点击率会是多少”,和线上真实曝光点击做对比。这样不切流量也能得到接近 AB 的效果。旁路日志要带上 requestId、用户行为序列、候选集、排序得分和最终曝光商品,缺一个字段都没法复盘。日志量比较大时要按天分桶存储,方便回溯。
6.2 延迟、覆盖率、状态体积三个硬指标
验证不能只看点击率。延迟看的是用户点击到特征生效的时延,对比 Kafka 消息里的 ts 和 Redis 特征更新时间,p95 压到 10 秒以内才算实时;覆盖率看有多少比例请求拿到了个性化候选,而不是全部落入热销兜底;状态体积看 Flink 作业的状态大小和 checkpoint 耗时,这两个指标直接决定长稳运行上限。三个硬指标都达标,再谈点击率提升才有意义。
我自己养成的习惯是,每个版本上线前先在测试环境完整跑一套旁路日志,连续观察三天趋势再决定要不要推全量。实时推荐链路很长,出问题往往不是模型的问题,而是埋点、状态、存储和网络里某个环节悄悄变了。这个教训让我现在遇到推荐效果波动时,第一反应永远是看 Kafka Lag 和 Redis 慢查询,而不是先去调排序模型参数。希望帮到你。
本文还有配套的精品资源,点击获取