简介:基于 Flink 的商品实时推荐系统完整项目源码,面向大数据、人工智能、物联网等专业的毕业设计、课程设计与 Flink 进阶学习者。系统以 Kafka 接收用户评分行为,由 Flink 完成实时与离线两类推荐:实时侧包括基于行为的推荐和实时热门统计,离线侧包括历史热门、历史优质商品与 ItemCF 推荐,覆盖从数据接入、流式处理、特征计算到结果落地的完整链路。压缩包共 408 个文件、约 4.27MB,核心代码以 44 个 Java 源文件为主,辅以 XML 配置、Vue/TS 前端页面、SQL 建表脚本和 properties 配置,目录结构清晰,便于运行调试和二次开发。项目已通过导师指导与答辩评审,代码经测试运行成功,附有详细文档和全部资料;目前已有 87 人学习下载,既适合在此基础上扩展推荐功能,也可直接用于毕设、课设或作业。
1. 商品实时推荐系统:先把这条链路上的角色对号入座
拿到这套基于 Flink 的商品实时推荐系统时,我第一反应是翻它的类文件清单。HbaseClient、ItemCFTask、OnlineRecommendMapFunction、HotProducts、HbaseSource、HbaseTableSource、TopNProductTask、StatisticsTask、HbaseSink——这些类名基本把架构画出来了:Kafka 负责送评分事件,Flink 负责实时和离线计算,HBase 负责存历史行为与推荐结果。实时推荐吃用户当下产生的评分行为,离线推荐吃历史聚合结果,两条链路最终都落到 HBase 里,供上层查询。这个项目适合正在做毕设、课设,或者想把 Flink + Kafka + HBase 整条链路搞明白的从业者,代码结构不复杂,但该有的环节都有。
2. Kafka 到 Flink 的评分流接入:数据结构、并行度与 offset 提交
2.1 评分事件的结构:user_id、item_id、score、timestamp 缺一不可
整套系统的数据源头是用户在 App 上的评分行为。Kafka 里每条消息对应的 JSON 结构大概是这样的:
{ "user_id": 1024, "item_id": 512, "score": 4.5, "timestamp": 1712841600000 }我在看这个项目的 HbaseClient 时注意到,它实际上是把行为记录去 HBase 的 user_behavior 表里做增量追加,而 Flink 消费 Kafka 后也会用同样的字段结构。也就是说,Kafka 里的消息格式和 HBase 里的列族字段必须保持一致。常见做法是在项目里单独建一个 RatingEvent POJO 类,用 Fastjson 或 Gson 做反序列化,字段名直接对应 JSON 的 key。
从接入角度看,timestamp 一定不要省。Flink 做事件时间处理时需要它来分配水位线,离线统计也需要它来界定“历史”的截止点。如果业务里没有这个字段,实时热门和离线 ItemCF 的时间窗口都会乱套。
2.2 FlinkKafkaConsumer 的并行度与消费位点:本地跑得通,集群上为什么乱
项目里接入 Kafka 的代码大致是:
Properties props = new Properties(); props.setProperty("bootstrap.servers", "localhost:9092"); props.setProperty("group.id", "rating-consumer-group"); props.setProperty("auto.offset.reset", "earliest"); FlinkKafkaConsumer<String> kafkaSource = new FlinkKafkaConsumer<>("product-rating", new SimpleStringSchema(), props); kafkaSource.setStartFromGroupOffsets(); DataStream<String> sourceStream = env.addSource(kafkaSource);这里有两个关键的参数。第一个是auto.offset.reset,设成earliest表示从最早可消费的位置开始,设成latest则只消费新消息。做推荐系统测试时,我一般用earliest,这样 replay 历史数据也能触发计算。第二个是group.id,它决定了消费者组的偏移量记录在 Kafka 的__consumer_offsets里,Flink 的成功恢复依赖这个组 ID。
在集群上最容易翻车的是并行度。Kafka 的 topic 如果分区数是 3,而 Flink 作业给 source 设置了并行度 8,那么只有 3 个并行子任务真正持有分区,另外 5 个空转。这本身不报错,但留给人的错觉是“作业很忙”。反过来,如果分区数大于并行度,同一时刻就有部分分区被轮询等待。实操中我会先把 topic 分区数定在并行度的整数倍,比如分区 12、并行度 6,让每个子任务稳定持有 2 个分区。源码里如果没写setParallelism,默认就是整个作业的并行度,这点很隐蔽。
2.3 HbaseSource 与 HbaseSink:rowkey 没设计好,一切白搭
项目里有 HbaseSource、HbaseTableSource、HbaseSink 三个类。HbaseSink 负责把 Flink 算出的推荐结果写入 HBase,HbaseSource 负责把历史行为读出来做离线计算,HbaseTableSource 则更像是给实时流做维表关联。它们的核心是 rowkey 设计。
我看到 CommonFlink 这类项目里的典型做法是:行为表 rowkey 拼成userId_reverseTimestamp_itemId,这样同一个用户的行为能按时间倒序连续存储,Scan 时指定 startRow 和 stopRow 就能取到最近 N 条。推荐结果表则直接以userId作为 rowkey,每行存多个推荐商品列,列名用rec_item_1、rec_item_2这种。
Put put = new Put(Bytes.toBytes(String.valueOf(userId))); put.addColumn(Bytes.toBytes("rec"), Bytes.toBytes("item_" + rank), Bytes.toBytes(String.valueOf(itemId))); hbaseSink.getTable().put(put);如果是往 HBase 里写,注意这里Bytes.toBytes的编码一致性。Java 的 String 默认 UTF-8,但如果 rowkey 里混了数字和字符串,建议统一拼成 String 再转字节,否则不同进程用String.valueOf(userId)拼出来的 rowkey 会不一致。我见过好几个直接把 userId 用Bytes.toBytes(userId)写、用Bytes.toString(bytes)读的案例,数字本身没问题,一旦拼接了分隔符,扫描范围立刻对不上。
3. 四类推荐结果的计算逻辑:ItemCF、实时行为、实时热门与历史统计
3.1 实时推荐:OnlineRecommendMapFunction 里到底做了什么
实时推荐的核心在 OnlineRecommendMapFunction。这类函数的典型逻辑是:每来一条用户评分事件,先把评分写入用户行为表,然后基于这个行为找到与该商品最相似的 TopK 商品,作为“看了 A 的人还看了 B”的实时输出。
为了拿到相似商品,它通常会去读 HBase 里预先算好的商品相似度矩阵。也就是说,离线 ItemCF 已经算好了 item 到 item 的相似度,实时推荐只是查询这张表。OnlineRecommendMapFunction 更像一个查询函数,而不是计算函数。它的内部结构大致是:
@Override public String map(RatingEvent event) throws Exception { String userId = event.getUserId(); String itemId = event.getItemId(); // 写入用户行为表,方便后续离线任务挖掘历史偏好 hbaseClient.incrementBehavior(userId, itemId, event.getScore()); // 从 item_similarity 表查与 itemId 最相似的 topN List<ItemScore> similarItems = hbaseClient.getTopSimilarItems(itemId, 10); // 过滤掉用户已经评分过的商品 List<String> history = hbaseClient.getUserHistory(userId); String recommendResult = similarItems.stream() .filter(rec -> !history.contains(rec.getItemId())) .map(ItemScore::getItemId) .collect(Collectors.joining(",")); return userId + ":" + recommendResult; }这段代码的逻辑说明:评分事件进来之后,先增量写行为,这一步是为了后续离线统计用的;然后从相似度表里取相似商品;最后把用户历史评分过的商品过滤掉,避免推荐重复内容。这里有个参数值得注意——getTopSimilarItems(itemId, 10)中的 10 是候选集大小,推荐系统里一般取 20 到 50 再过滤,最后只剩 3 到 5 个。如果一开始就取 10 个,过滤完后可能只剩 2 个,输出列表显得很空。这个数值建议直接调大一点,缓存压力也不大。
3.2 实时热门:滑动窗口 + 热度衰减的热门榜
实时热门在 HotProducts 类里完成。这部分的实现思路是用 Flink 的滑动窗口统计一定时间窗口内的商品评分次数,再按次数排序输出热门商品。常见做法是:
DataStream<RatingEvent> input = env.addSource(kafkaSource); DataStream<ProductCount> hotStream = input .keyBy(e -> e.getItemId()) .window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1))) .aggregate(new CountAggregate(), new WindowResultFunction());这里SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1))的含义是:窗口长度 10 分钟,滑动步长 1 分钟,即每 1 分钟输出一次最近 10 分钟的热门榜单。参数调整要看业务节奏,电商大促期间窗口可以缩短到 5 分钟,普通场景 10 分钟比较稳。如果数据量大,窗口内状态会比较大,建议把window的触发器改成CountTrigger,比如每 5000 条数据触发一次输出。
热度计算不能只用简单的次数,否则老商品永远占着榜单。项目里如果要做得更自然一点,可以在CountAggregate里给评分加上时间衰减因子,比如weight = score * exp(-0.5 * ageInMinutes)。但我在源码里看到的是纯次数统计,这更符合课程设计的定位。你要是想更真实,可以在窗口函数里把最新时间戳和最早时间戳做差值,对计数做衰减。
3.3 离线 ItemCF:共现矩阵与相似度计算
ItemCF 的离线计算在 ItemCFTask 里。ItemCF 的核心是:如果用户同时给商品 A 和商品 B 打了分,那 A 和 B 就算一次共现。所有用户打分行为汇总后,得到共现矩阵,再除以商品的热度得到相似度。
具体流程一般分三步:
第一步,从 HBase 的 user_behavior 表读出全部历史评分记录。这里我建议用 HbaseSource 做全量 Scan,而不是只读最近几天。离线任务不追求秒级,全量数据算出的相似度更稳定。
第二步,用 Flink 的 DataSet 或 DataStream 按 userId 做 groupBy,在同一个用户的行为集合内部生成商品对:
DataStream<RatingEvent> history = env.addSource(hbaseSource); DataStream<Tuple2<String, String>> pairs = history .keyBy(e -> e.getUserId()) .window(TumblingProcessingTimeWindows.of(Time.days(1))) .process(new GenerateItemPairs());生成商品对时要注意去重。一个用户如果在一个窗口内给同一件商品打了三次分,只应该算一次共现,否则相似度会被单个用户放大。处理方式是把商品列表转成 TreeSet,再去生成不重复的排列组合。
第三步,统计共现次数并计算相似度。相似度常用余弦公式:
double sim = cooccurrence.get(countAAndB) / Math.sqrt(countA * countB);然后把结果写入 HBase 的 item_similarity 表,rowkey 用itemA_itemB,列族存相似度值。这个表的规模是 N 乘 N,N 是商品数。如果商品数上万,全量写入会产生大量 rowkey,建议把相似度小于 0.01 的过滤掉,只在表中保留有价值的边。
3.4 StatisticsTask 和 TopNProductTask:离线统计的两条分支
StatisticsTask 主要负责历史热门商品和历史优质商品的统计。历史热门就是按历史评分次数排序,历史优质商品则往往要考虑评分均值。如果只有次数没有均值,一款 1 星但被刷了 1000 次的商品会进榜单,所以更合理的统计是:
DataStream<ProductScore> stats = history .keyBy(e -> e.getItemId()) .reduce(new ReduceFunction<RatingEvent>() { @Override public RatingEvent reduce(RatingEvent e1, RatingEvent e2) { double newScore = e1.getScore() + e2.getScore(); long newCount = e1.getCount() + e2.getCount(); return new RatingEvent(e1.getItemId(), newScore / newCount, newCount); } });这里我用了reduce而不是aggregate,是因为 reduce 可以把均值状态直接塞进事件结构里,减少自定义 Accumulator。注意newScore / newCount在整数除法下会丢掉小数,先把 score 转成 double 再除。这个坑很常见。
TopNProductTask 则是把各类统计结果汇总排序,取出全局 TopN。它的数据来源是前面算好的多个中间结果表,合并后统一排序。这里要考虑的是合并时的权重:历史热门、历史优质、ItemCF 三张表的分数量纲不同,不能直接相加。常见做法是各表先归一化到 0 到 1,再加权求和。权重参数一般放在配置里,避免硬编码。
4. 避坑指南:Flink + Kafka + HBase 联动最常见的五个翻车点
4.1 现象:HbaseSink 写入频繁报 RegionTooBusy
原因:写入请求的 rowkey 设计成随机字符串,导致写入压力分布在整个 Region 的所有节点上。如果 rowkey 以 UUID 开头,HBase 的写请求会散落到各个 RegionServer,表面看起来均衡了,但每次批量提交的 Put 对应多个 Region,很容易把 RegionServer 的队列打满。解决:把 rowkey 前缀改成用户分片,比如userId % 100固定为两位前缀,让同一个用户的数据落在一个 Region 里。同时把 HbaseSink 的BufferedMutator缓冲大小调到 4MB 到 8MB,减少 RPC 次数。
4.2 现象:水位线不动,消费速率上不去
原因:Kafka Topic 分区数远小于 Flink 作业并行度,且数据源里的setStartFromGroupOffsets()在没有新消息时永远停在当前位点。另一个常见原因是事件时间字段没有正确提取,水位线一直停留在初始时间,窗口永远不触发。解决:先确认 Kafka 分区数,并行度不超过分区数;然后在 Flink 代码里显式调用assignTimestampsAndWatermarks,用WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(10))加上timestamp字段提取器。如果仍然不动,看 Flink UI 上每个 Subtask 的 Input Watermark 数值,卡在最小值说明 source 没拿到数据。
4.3 现象:ItemCF 推荐出来的全是历史热门
原因:ItemCF 的相似度矩阵没有做热门商品惩罚。比如一个热门商品 A 和千万个商品都共现过,余弦相似度分母虽然用了平方根,但热门商品本身的出现次数会让很多商品都跟它相似。解决:在计算相似度时加入流行度惩罚,常用方案是把countA * countB的分母改成Math.pow(countA, 0.5) * Math.pow(countB, 0.5)改为Math.pow(countA * countB, 0.6),放大热门商品的惩罚。或者直接过滤掉出现次数排名前 0.1% 的商品。
4.4 现象:HbaseTableSource 查询刚写入的数据查不到
原因:HBase 的写入是先写 MemStore,再落 HFile,如果你的 Flink 查询 task 在写入任务还没 flush 时就去读,会看到不一致。特别是实时链路里,HbaseSink 刚完成put,下一层立刻用 HbaseTableSource 去查,容易扑空。解决:开启 HBase 的hbase.client.scanner.timeout.period和设置setCaching(100)不一定能解决数据可见性问题。更直接的办法是让后续查询延迟几毫秒,或者在 HbaseSink 里强制table.flushCommits()后再返回。如果是流式关联,用AsyncFunction配合重试三次,每次间隔 100ms,比调整 HBase 参数更靠谱。
4.5 现象:重启作业后用户立刻被重复推荐
原因:Flink 的 Checkpoint 里保存了 Kafka 消费位点和 HBase 写入的中间状态,但 HBase 里的推荐结果表没有做幂等。重启后作业从 Checkpoint 恢复,可能重新处理一部分消息,把相同的推荐结果再次写入 HBase。解决:给 HBase 的推荐结果表加一个带时间戳的列,比如rec_ts,写入时先检查该 rowkey 上rec_ts是否大于当前事件时间,大于则跳过。更简单的做法是把推荐结果表的 TTL 设置为窗口长度,比如 10 分钟,这样即使重复写入,旧数据也会被自动过期。
5. 端到端跑通:从建表、灌模拟数据到验证推荐结果
5.1 准备环境:HBase 表设计、Kafka Topic 和 Flink 作业打包
先把 HBase 需要用的两张表建出来。行为表和推荐结果表:
create 'user_behavior', 'info', 'action' create 'item_similarity', 'sim' create 'recommend_result', 'rec'这里user_behavior的列族action存用户的评分行为,每个用户一行,列名是itemId_score_timestamp的形式。item_similarity的sim列族存相似度矩阵。recommend_result的rec列族存实时与离线的推荐结果。
Kafka 这边创建 topic:
kafka-topics.sh --create --zookeeper localhost:2181 \ --replication-factor 1 --partitions 6 --topic product-ratingTopic 的 partition 设置为 6,是为了和后面 Flink 作业的并行度对齐。如果你只有一台机器,replication-factor可以设 1,否则 Kafka 会报警告。
Flink 作业打包用 Maven:
mvn clean package -DskipTests打包前检查pom.xml里的flink-shaded-hadoop依赖是否注释掉。如果集群上已经有 HBase 客户端,本地打包时就不要把 HBase 的 jar 打进去,否则 ClassCastException 会教你做人。我在项目文档里看到参考资料建议用provided作用域,这点很实用。
5.2 模拟评分流:用一段脚本持续灌数据
没有真实业务流量时,我们用 Python 脚本模拟用户评分。脚本每秒随机生成一个用户 ID 和商品 ID,发送到 Kafka:
import json import random import time from kafka import KafkaProducer producer = KafkaProducer( bootstrap_servers='localhost:9092', value_serializer=lambda v: json.dumps(v).encode('utf-8') ) while True: event = { 'user_id': random.randint(1, 1000), 'item_id': random.randint(1, 500), 'score': round(random.uniform(1.0, 5.0), 1), 'timestamp': int(time.time() * 1000) } producer.send('product-rating', value=event) time.sleep(1)脚本里的参数说明:random.randint(1, 1000)控制用户规模,item_id控制在 1 到 500,这样 ItemCF 的共现矩阵不会太大。time.sleep(1)代表每秒一条消息,实际压测时可以改成 0.1 秒甚至 0.01 秒,看集群吞吐。如果想让某个用户的行为更集中,可以把user_id固定为几个值,这样更容易在离线结果里看到某个用户的偏好变化。
5.3 验证点:实时推荐是否被新行为影响、离线结果是否收敛、热门榜是否滚动
启动 Flink 作业后,先订阅 Kafka 控制台确认消息到达:
kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic product-rating --from-beginning --max-messages 10然后观察三个验证点。
第一个验证点是实时推荐的灵敏度。用同一个user_id持续给一个冷门商品打分,然后去 HBase 查这个用户的recommend_result表,如果 rowkey 对应的推荐商品列表发生了变化,说明实时链路在起作用。如果始终不变,检查 OnlineRecommendMapFunction 里是不是忘记把新行为写入行为表了。
第二个验证点是离线 ItemCF 是否收敛。跑完离线任务后,随机抽查一个商品的相似商品列表,看它是否跟该用户历史行为有关。比如用户只看过科幻电影,推荐结果里突然出现言情片,而且相似度高于 0.8,说明共现矩阵可能有问题。
第三个验证点是热门榜是否滚动。持续运行 15 分钟,每 1 分钟记录一次 TopN 热门商品。如果榜单完全不变,说明滑动窗口的事件时间没有正确推进;如果榜单剧烈抖动,说明窗口太短或数据量不均匀。正常情况应该是前几名缓慢变化,中间名次有小幅波动。
6. 从 HBase 把推荐结果读出去的三种姿势与最终排重
6.1 直接 Scan:适合 demo,但要注意 rows 限制
最简单的方式是查recommend_result表,用 Scan 指定 rowkey 前缀:
Scan scan = new Scan(); scan.setStartRow(Bytes.toBytes("100001|")); scan.setStopRow(Bytes.toBytes("100001|")); scan.setCaching(20);这里setCaching(20)控制每次 RPC 拉取多少行,值过小会很慢,过大则容易超时。Demo 场景够用,但并发一高,Scan 会把 RegionServer 的资源吃满。所以只能是原型阶段的过渡方案。
6.2 HbaseTableSource 流式关联:小心查询超时
如果要在实时推荐返回结果时同时查 HBase 里的用户历史,建议用HbaseTableSource配合AsyncFunction:
AsyncDataStream.orderedWait( inputStream, new HBaseAsyncLookupFunction(tableName), 3, TimeUnit.SECONDS, 5);3是超时时间,5是最大并发请求数。超时时间设得过短,高峰期查询就大量失败;设得过长,背压会传导到 Kafka 消费端。我一般先测 HBase 单次 Get 的 P99 耗时,再设置成 P99 的两倍作为超时阈值。
6.3 每个用户一张 TopN 行:rowkey 顺序决定了读写效率
最终对外提供查询时,我会把每个用户的 TopN 推荐结果落到一行里,rowkey 拼接顺序是userId 倒序时间戳。这样用户刚产生的推荐结果永远在行的最前面,后续查询时最新的数据无需跳跃。排重的技巧是写之前先读一次旧列表,把新列表合并时去掉重复 itemId,只保留分数最高的那条。
我从这个项目里学到的习惯是:任何从 HBase 读出来的列表,都要在业务层再做一次去重和过滤,因为 HBase 不保证同一个 rowkey 下多个列的写入顺序,也不保证上次覆盖一定能立即生效。从那以后我每次跑推荐任务都会强制走一遍“写入前读旧、写入后校验”的流程,再配合 TTL 让过期数据自动消失。这套流程虽然麻烦,但至少不会在半夜接到线上数据重复的告警。希望帮到你。
本文还有配套的精品资源,点击获取