简介基于 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); FlinkKafkaConsumerString kafkaSource new FlinkKafkaConsumer(product-rating, new SimpleStringSchema(), props); kafkaSource.setStartFromGroupOffsets(); DataStreamString 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 与 HbaseSinkrowkey 没设计好一切白搭项目里有 HbaseSource、HbaseTableSource、HbaseSink 三个类。HbaseSink 负责把 Flink 算出的推荐结果写入 HBaseHbaseSource 负责把历史行为读出来做离线计算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 ListItemScore similarItems hbaseClient.getTopSimilarItems(itemId, 10); // 过滤掉用户已经评分过的商品 ListString 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 的滑动窗口统计一定时间窗口内的商品评分次数再按次数排序输出热门商品。常见做法是DataStreamRatingEvent input env.addSource(kafkaSource); DataStreamProductCount 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在同一个用户的行为集合内部生成商品对DataStreamRatingEvent history env.addSource(hbaseSource); DataStreamTuple2String, 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 乘 NN 是商品数。如果商品数上万全量写入会产生大量 rowkey建议把相似度小于 0.01 的过滤掉只在表中保留有价值的边。3.4 StatisticsTask 和 TopNProductTask离线统计的两条分支StatisticsTask 主要负责历史热门商品和历史优质商品的统计。历史热门就是按历史评分次数排序历史优质商品则往往要考虑评分均值。如果只有次数没有均值一款 1 星但被刷了 1000 次的商品会进榜单所以更合理的统计是DataStreamProductScore stats history .keyBy(e - e.getItemId()) .reduce(new ReduceFunctionRatingEvent() { 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 这边创建 topickafka-topics.sh --create --zookeeper localhost:2181 \ --replication-factor 1 --partitions 6 --topic product-ratingTopic 的 partition 设置为 6是为了和后面 Flink 作业的并行度对齐。如果你只有一台机器replication-factor可以设 1否则 Kafka 会报警告。Flink 作业打包用 Mavenmvn clean package -DskipTests打包前检查pom.xml里的flink-shaded-hadoop依赖是否注释掉。如果集群上已经有 HBase 客户端本地打包时就不要把 HBase 的 jar 打进去否则 ClassCastException 会教你做人。我在项目文档里看到参考资料建议用provided作用域这点很实用。5.2 模拟评分流用一段脚本持续灌数据没有真实业务流量时我们用 Python 脚本模拟用户评分。脚本每秒随机生成一个用户 ID 和商品 ID发送到 Kafkaimport json import random import time from kafka import KafkaProducer producer KafkaProducer( bootstrap_serverslocalhost:9092, value_serializerlambda 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, valueevent) 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配合AsyncFunctionAsyncDataStream.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 让过期数据自动消失。这套流程虽然麻烦但至少不会在半夜接到线上数据重复的告警。希望帮到你。本文还有配套的精品资源点击获取