简介:一套基于Flink的电商商品实时推荐系统项目资料,面向大数据方向在校生、毕业设计开发者以及Flink初学者,意在解决用户评分行为驱动下的实时与离线推荐问题。项目通过Kafka接收评分数据,由Flink完成实时推荐和离线推荐:实时侧包括基于行为推荐与实时热门统计,离线侧包括历史热门、历史优质商品和ItemCF协同过滤,同时借助HBase完成特征存储、结果写入与源表读取,形成完整的数据闭环。资源共408个文件,主要类型包含44个Java核心源码、245个XML配置、Vue与TypeScript前端工程、SQL建表脚本及CSV样本数据,配套文档与工程齐全,压缩包仅4.27MB,结构紧凑、模块清晰,便于快速部署与后续扩展。已有86人学习下载。项目中ItemCFTask、StatisticsTask、HotProducts、TopNProductTask、OnlineRecommendMapFunction等模块,清晰展示了离线协同过滤、统计聚合、热门排行与在线推荐映射的实现思路,代码经运行验证,可直接用于课程设计、毕业设计或项目答辩,也适合在此基础上深入阅读和二次开发。
1. 从一次评分到推荐结果:这套链路为什么值得被讲透
用户在商品页点了五颗星,或者打了四分的评价,这个动作在传统架构里只会落进数据库,等凌晨的批处理任务去算一次热门榜。但放到推荐场景里,这一个评分其实是用户当下兴趣最强烈的信号——他刚看完商品详情、比过价格、翻过评论,最后愿意花两秒钟打分,说明此刻的偏好是明确的。如果系统能在他停留的这几秒里把这个信号用起来,推荐位的点击率往往比只看历史画像高出一截。
标题里的这条链路,本质上是把「行为事件流」和「用户历史画像」两条数据流在 Flink 里做一次融合:Kafka 负责把评分事件实时送进来,Flink 一边用滑动窗口捕捉用户最近几分钟的短期兴趣,另一边把用户长期的历史评分读进内存模型做协同过滤,两条结果合并后才去召回候选商品。适合谁呢?适合那些已经有用户行为埋点、却还停留在每晚离线算推荐的团队;也适合想弄清楚 Flink 在推荐系统里到底该承担「实时计算」还是「全量训练」这两种角色的读者。整套方案不依赖特定的机器学习平台,Flink 集群加上 Kafka 就能落地,关键是把状态、窗口和维表 Join 这三个基本功用到位。
2. 推荐系统的数据管道设计:Kafka Topic 拆分与 Flink 接入方式
2.1 评分事件的 Topic 划分原则
Kafka 里的 Topic 设计直接决定 Flink 作业的拓扑复杂度。很多刚接触实时推荐的团队会把所有用户行为塞进一个叫user_behavior的 Topic,然后用一大段CASE WHEN在 Flink 里分流。这种做法在数据量小的时候看不出问题,一旦评分、浏览、加购、收藏四类事件的量级同时涨起来,单 Topic 的多消费者组会互相拖慢,Flink 作业的反压监控也很难定位是哪一类事件导致的。
常见的做法是按事件类型拆分 Topic:rating_events、view_events、cart_events、purchase_events,分区数按事件量的预估峰值除以单分区吞吐上限来定。评分事件一般只有浏览量的十分之一到二十分之一,分区数设成 6 到 12 就够;浏览事件如果日活百万级别,分区数至少 24 起。分区的意义不只是吞吐,它还决定了 Flink 的并行度上限——一个 Flink 算子实例最多对应一个 Kafka 分区。
2.1.1 消息键的选择与用户维度的数据局部性
评分消息写入 Kafka 时的 key 建议直接用userId,不要用随机字符串。原因在于 Flink 消费后要做keyBy(userId)才能把同一个用户的评分聚到同一个算子实例上。如果 Kafka 端 key 和 Flink 端 keyBy 不一致,就会出现跨实例的数据重排,网络开销和序列化开销同时上升。Kafka 生产者端的 key 策略是userId.toString(),这样同一个用户的所有评分天然落在同一个分区,Flink 消费端即使不做rebalance,也能大概率在本地完成后续的窗口聚合。
2.2 JSON 序列化方案的取舍
评分消息的 payload 一般会包含userId、itemId、score、timestamp、sceneId五个字段。序列化格式推荐使用 JSON,虽然它比 Avro 多出 20% 到 30% 的体积,但对于日千万级事件量来说,这点体积换来的排查便利性很值。Flink 里用JSONDeserializationSchema接 Kafka,把原始字符串解析成RatingEventPOJO,代码里最需要注意的是时间戳字段的处理。
public class RatingEvent { public long userId; public long itemId; public double score; public long timestamp; public int sceneId; public static RatingEvent fromJson(String json) throws IOException { ObjectMapper mapper = new ObjectMapper(); JsonNode node = mapper.readTree(json); RatingEvent event = new RatingEvent(); event.userId = node.get("userId").asLong(); event.itemId = node.get("itemId").asLong(); event.score = node.get("score").asDouble(); event.timestamp = node.get("timestamp").asLong(); event.sceneId = node.has("sceneId") ? node.get("sceneId").asInt() : 0; return event; } }这段解析逻辑里有两个容易被忽略的细节:sceneId是可选字段,线上历史数据可能没有这个字段,解析时必须用has()判断,否则一条脏数据就会让整个作业的fromJson抛异常;timestamp字段不要用System.currentTimeMillis()在 Flink 端补,因为 Kafka 生产者所在的应用服务器和 Flink 集群之间可能有毫秒级的时间偏差,对于窗口计算来说,采用事件自带的时间戳才准确。
2.3 Flink SQL 还是 DataStream API
评分事件接入 Flink 后,接下去是路由选择的问题:用 Flink SQL 还是 DataStream API。标题里的场景同时涉及窗口聚合和实时召回计算,建议是这两者的混用——用 Flink SQL 做清洗、过滤、去重这些相对标准的操作,用 DataStream API 做需要深挖状态或自定义触发逻辑的部分。
举例来说,用户在一个 session 内可能对同一个商品评分多次,只保留最后一次评分这个动作,用 SQL 写需要开窗加ROW_NUMBER(),而 Flink 的KeyedProcessFunction里直接维护一个ValueState<Long>存上次评分时间,两条语句就能解决。SQL 的优点是开发速度快,DataStream 的优点是状态控制灵活,两者通过TableEnvironment.toDataStream()和StreamTableEnvironment.fromDataStream()互相转换,作业内部不会产生额外的序列化开销。
3. 实时推荐的算法核心:基于用户行为的评分预测与物品召回
3.1 短期行为权重与评分归一化
实时推荐的「实时」二字,主要体现在用户最近几分钟的行为对推荐结果的即时影响上。用户过去三十天的平均评分可能是 3.8 分,但最近十分钟他连续给三本书打了五星,说明他当下的阅读兴趣正在往某个方向倾斜。如果只用历史均值,这个信号会被稀释掉。
需要设计一套权重公式:评分事件的权重按时间衰减,以Math.exp(-elapsedMinutes / 30.0)作为衰减因子,30 是半衰期参数,表示 30 分钟前的评分对当前兴趣的影响只有刚发生时刻的1/e,约 37%。再把评分值归一化到 0 到 1 的区间,公式为normalizedScore = (rawScore - 1.0) / 4.0,这样处理的原因是原始的 1 到 5 分制里,3 分和 4 分之间的差异与 1 分和 2 分之间的差异在实际偏好强度上并不等价,归一化后参与相似度计算更稳定。
DataStream<ItemScore> weightedScores = ratingStream .keyBy(event -> event.userId) .process(new KeyedProcessFunction<Long, RatingEvent, ItemScore>() { private ValueState<Double> userAvgState; private ValueState<Long> lastUpdateState; @Override public void processElement(RatingEvent event, Context ctx, Collector<ItemScore> out) throws Exception { double avgScore = userAvgState.value() == null ? 3.0 : userAvgState.value(); double weight = Math.exp(-elapsedMinutes(event.timestamp) / 30.0); double normalized = (event.score - avgScore) / 4.0 + 0.5; double score = normalized * weight; ItemScore itemScore = new ItemScore(); itemScore.userId = event.userId; itemScore.itemId = event.itemId; itemScore.score = score; itemScore.timestamp = ctx.timerService().currentProcessingTime(); out.collect(itemScore); userAvgState.update(avgScore * 0.95 + event.score * 0.05); } });代码里做了两个关键设计:normalized的计算不是简单地(rawScore - 1.0) / 4.0,而是减去了该用户的历史平均分,这叫 User-Centric Normalization,可以消除不同用户打分尺度的差异——有人习惯打 2 到 3 分,有人习惯打 4 到 5 分,减去均值后,同样的原始分数变化在不同用户间就变得可比了;userAvgState.update(avgScore * 0.95 + event.score * 0.05)是一个滑动平均,用 5% 的学习率让用户均值缓慢漂移,避免单次极端评分瞬间拉偏整体均值。
3.2 实时协同过滤:Hash-based 最近邻召回
实时推荐阶段不可能跑全局的 ALS 矩阵分解,因为矩阵分解的迭代训练耗时以分钟计,等模型算完,用户当前的兴趣窗口已经过去了。业界的常规做法是用一个近似的最近邻召回:把用户近期高权重的评分商品作为种子,在商品相似度矩阵里查 Top N 相似商品。这个商品相似度矩阵是离线算好的,存放在 Redis 里,Flink 作业在运行期用 Async I/O 去查询,而不是把相似度矩阵也塞进 Flink 状态——矩阵可能几十万乘几十万,全放状态里内存吃不消。
Redis 里相似度矩阵的 key 设计为item_sim:{itemId},value 是类似itemId1:0.87,itemId2:0.76,itemId3:0.65的字符串,Flink 端用RedisAsyncLookupFunction批量获取候选商品。查询的种子商品数量不要太多,取用户最近 20 个不同商品的评分中权重最高的 5 个即可。5 个种子商品每个取 20 个相似商品,候选池在一百个商品左右,这个量级足够应付精排阶段的排序了。
3.2.1 候选商品的去重与过滤
召回到的候选商品不能直接输出,要经过一层过滤:用户已经打过分且分数高于 3 的商品直接排除,因为推荐位放一个用户明确评价过的东西没有转化意义;商品本身有上下架状态,下架商品要在 Redis 里维护一个blacklist的 Set,Flink 每五分钟从这个 Set 拉一次增量更新放进 BroadcastState,过滤时查这个状态比每次查 Redis 省掉大量网络开销。
3.3 实时与离线的结果融合排序
实时召回结果和离线推荐结果不能简单地按实时优先排列,因为实时召回只有五六个种子,覆盖面窄,容易让用户看到全是同类型商品。常见融合策略是:实时召回的候选排在最前,但同一品类不超过三个,剩下的位置从离线推荐结果里补;两种来源的候选共用一个 CTR 预估分,评分格式统一后就可以混合排序。
这里必须处理一个问题:实时部分给的分数和离线部分给的分数不在同一个量纲上。实时分数有时间的衰减因子,离线分数是ALS预测值加规则加成,直接相加等于让离线分主导。实际操作时对两个分数分别做 Min-Max 归一化到 0 到 1 的区间,再加权求和,权重系数0.65给实时、0.35给离线,这个比例在绝大多数电商场景下比五五开表现好,因为实时信号虽然强,但稀疏,占比过高会让排序结果抖动得非常厉害。
4. 离线推荐的批流一体实现:ALS 模型训练与周期性更新策略
4.1 离线数据源的抽取方案
离线推荐需要全量用户历史评分数据,这些数据存在业务库的user_rating表和 Kafka 里消费过的历史事件中。如果 Kafka 的留存时间只有三天,那三天前的评分就全丢了,因此离线数据的基础来源应该是数据库或者数据仓库。常用的做法是直接用 Flink CDC 把 MySQL 里的user_rating表全量同步到 Hive 表,再通过 Flink SQL 的批模式读取 Hive 表做模型训练。热词里提到的 mysql增量同步工具选型,在这个场景下的推荐是 Flink CDC 本身——它不需要额外部署独立同步进程,对这张表的 binlog 实时监听,每天凌晨定时把全量快照写进 Hive 分区就够了。
从 Kafka 消费的历史数据也可以作为补充,比如用户浏览行为比评分行为丰富得多,但浏览行为不在这张 MySQL 表里。处理方式是把 Kafka 里的浏览事件通过INSERT INTO hive_table SELECT ...定期落成 Hive 分区表,然后训练脚本用INSERT OVERWRITE把两个数据源做UNION ALL合并。注意评分数据和浏览数据在训练样本里的权重需要调节,浏览一个商品只算是弱正样本,评分 4 分以上才是强正样本,样本权重分别设为 0.3 和 1.0 比较合理。
4.2 Flink 批任务训练 ALS 模型
Flink 的批处理能力在这个场景里主要体现在训练数据的预处理上——清洗、过滤、用户商品交叉过滤,这些都是典型的批任务。真正跑 ALS 矩阵分解的算法可以放在 Flink ML 库里,也可以把预处理后的三元组数据输出成一个文本文件,交给 Spark MLlib 或者本地 Python 脚本训练。后一种做法更灵活,因为 Flink ML 的 ALS 实现和 Spark 相比,更新频率低、示例少,踩坑时排查成本高。
如果保留在 Flink 批任务里做预处理,核心代码是一段 SQL:
CREATE TABLE rating_train_data AS SELECT user_id, item_id, score FROM ( SELECT user_id, item_id, score, ROW_NUMBER() OVER (PARTITION BY user_id, item_id ORDER BY rating_time DESC) AS rn FROM rating_source ) t WHERE t.rn = 1 AND score >= 2.0 AND user_id IN (SELECT user_id FROM active_users WHERE active_days >= 7);这个预处理干了三件事:ROW_NUMBER()去重,确保每个用户对每个商品只保留最新一条评分;分数小于 2 的负向评价直接过滤掉,因为负向评分对协同过滤的训练有干扰,用户不会因为有商品是他讨厌的就会喜欢它的近似商品;只保留近七天活跃用户,把那些注册后从没回来过的僵尸用户从训练集里剔除,否则 ALS 的隐因子空间会被大量空行拖慢收敛。预处理产出的三元组数据量级如果超过千万行,ALS 的rank参数设置在 20 到 50 之间、iterations在 10 到 20 之间是合理范围;隐因子维度设太高容易过拟合,设太低表示不了复杂的用户偏好结构,20 是冷启动场景的常见起点。
4.3 离线圈的调度周期与碰撞规避
离线模型训练任务和实时作业必须在物理或逻辑上隔离,这是容易踩的一个大坑。如果离线圈和实时圈跑在同一个 Flink 会话集群上,凌晨两点的 ALS 训练任务会把 TaskManager 的 CPU 和内存吃满,导致实时推荐作业在凌晨出现长达几十分钟的延迟高峰——恰好是用户活跃度另一个小高峰的时段。
常见的落地方式是分成两套 Flink 集群:一套专职跑流式作业,用yarn-session模式常驻,资源固定;另一套用yarn-per-job模式跑批任务,任务结束后资源自动释放。调度周期方面,离线模型每六小时训练一次比较合理:太频繁则资源和时间成本高,太稀疏则推荐结果跟不上新商品的入库节奏。如果商品池每天新增几千个商品,离线模型最好增加「近 24 小时新品补充召回」的逻辑,否则新品在六小时窗口内没有任何评分历史,协同过滤模型会直接把它们漏掉。
5. Flink 作业的关键参数调优与 Kafka 数据一致性保障
5.1 Checkpoint 配置与状态后端选型
实时推荐作业的状态主要是两类:每个用户的滑动平均分和窗口聚合的中间结果,状态量级不大,单个用户几十字节,百万用户也就几十 MB。基于这个体量,状态后端直接用 RocksDB 是浪费的——RocksDB 适合的是单 Key 大 Value 或者状态总量超过 TaskManager 堆内存的场景。几十 MB 的配置用 Heap StateBackend 就足够了,访问速度快,没有序列化开销。
Checkpoint 的间隔建议设成 30 秒,比默认的 60 秒短一点。原因是推荐作业对延迟敏感,如果中间 60 秒的状态没做 checkpoint,这期间计算出的所有用户权重都会在故障恢复时丢失。30 秒的代价是 checkpoint 更频繁地做增量快照,对这个小状态体量的作业来说几乎感觉不到。minPauseBetweenCheckpoints同样设成 30 秒,避免 checkpoint 之间腾不出间隔,导致前一次还没完成、后一次已经开始,那样反而不稳定。
execution.checkpointing.interval: 30s execution.checkpointing.min-pause: 30s execution.checkpointing.tolerable-failed-checkpoints: 3 state.backend.type: heap state.checkpoint-storage: filesystem state.checkpoints.dir: hdfs:///flink-checkpointstolerable-failed-checkpoints这个参数在推荐场景里值得单独拿来讨论。默认值是 0,意味着连续两次 checkpoint 失败作业就挂了;但推荐作业的价值是持续输出,短暂的几个 checkpoint 失败还不至于让推荐结果错到不可接受,设成 3 可以有效避免因为一次网络抖动导致整个作业重启,尤其是 Kafka 和 Flink 之间的连接出现瞬时异常的状况。
5.2 Kafka 消息延迟高与数据重复的排查视角
热词里 kafka消息延迟高和数据重复这两个词,在推荐场景中几乎必然会碰到。延迟高的根因往往不是 Kafka 本身,而是 Flink 的消费速度跟不上生产速度。优先查看 Flink Web UI 里各算子实例的backPressure指标,如果 Source 算子的backPressure比例超过 40%,说明下游处理能力是瓶颈。排查步骤是先看keyBy之后有没有数据倾斜——某个userId的评分量占到总量的百分之二三十时,分配到这个 Key 对应的算子实例自然被拖慢。
数据重复的现象则更隐蔽:Kafka 的生产者重试机制在分布式环境下可能导致消息被写入多次,消费者的enable.auto.commit如果设成false但 Flink 端没有正确维护 checkpoint 里的 offset,恢复重放时就会出现重复消费。处理办法是在 Flink 的消息解析阶段加入去重逻辑——核心是基于userId + itemId + timestamp这个三元组,在 KeyedProcessFunction 里维护一个ValueState<Long>记录上次处理到的 timestamp,如果新到的消息 timestamp 不大于状态里的值,直接丢弃。这个去重策略能挡住绝大多数重复消息。另外要强调的是,Flink 的 exactly-once 保证依赖 Kafka Source 的setStartFromLatest或setStartFromEarliest配合 checkpoint 机制,不要试图通过enable.auto.commit手动控制 offset 与 Flink 的 checkpoint 共存,两者会互相冲突。
5.3 实时推荐结果的输出与链路验证
推荐结果算完之后的输出方式也很关键。常见的做法是把推荐结果写回 Redis,key 为rec:{userId}:{sceneId},value 是商品 ID 列表的 JSON 序列化,TTL 设成 15 分钟。为什么是 15 分钟而不是更长?因为推荐结果要跟着用户的实时行为走,写死太长就失去了实时的意义;太短则会频繁触发 Flink 的周期输出,增加无谓的吞吐压力。同时把用户的最新评分行为写一份到 Hive 分区表realtime_behavior_log,供后续版本迭代模型时回溯分析。
5.3.1 onsumer 消费位置的校准技巧
上线验证阶段有个技巧:先不要直接在生产消费最新数据,把这个实时推荐作业连接到一个独立的 Kafka 消费组,从最近 24 小时的 offset 开始回放,看这段历史行为数据上跑出来的推荐结果是否合理。比如找几个有明确长短期偏好差异的用户,验证他们最近一小时的实时推荐列表是否明显偏向了新兴趣点,同时验证离线补充的商品又没有完全被实时结果挤掉。这个验证方法避免了新作业一上线就处在「从当前 offset 开始计算但窗口里还没攒够数据」的空窗期——这个阶段推荐结果为空,线上会以为作业挂掉了。
5.3.2 链路延迟的监控阈值
最后提一个监控阈值:从 Kafka 收到评分消息到 Redis 里出现更新后的推荐结果,端到端延迟正常应该控制在 3 到 8 秒。如果延迟超过 15 秒,优先看 Redis 的读写耗时——这类场景 Redis 通常会打到一个 2 到 5 毫秒的延迟,但如果Async I/O的超时参数设得太小导致大量请求被丢弃,延迟就会直接飙升。这个链路延迟建议直接对接进已有的 Kafka 监控大盘,比如用 Kafka 的consumer_lag指标结合 Flink 作业的numRecordsInPerSecond做关联,哪一侧出现倾斜就能快速定位到具体环节。
本文还有配套的精品资源,点击获取