news 2026/9/26 7:29:17

Flink+Kafka+HBase商品实时推荐系统实战:源码解析与避坑指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink+Kafka+HBase商品实时推荐系统实战:源码解析与避坑指南

简介:基于 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-rating

Topic 的 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 让过期数据自动消失。这套流程虽然麻烦,但至少不会在半夜接到线上数据重复的告警。希望帮到你。

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/26 7:29:00

2026专科生论文降AI率工具测评:从原理到实操的完整指南

2026专科生毕业季最让人抓狂的事&#xff0c;不是论文写不出来&#xff0c;而是写完了、查重过了&#xff0c;结果卡在学校新加的一道门槛上——AI率检测。前阵子陪表弟改论文&#xff0c;他第一稿用AI工具搭了个大概&#xff0c;自己润色了一部分&#xff0c;结果学校系统一查…

作者头像 李华
网站建设 2026/9/26 7:28:48

Mac M系列芯片部署Qwen-Image-Lightning全栈指南

1. 项目概述&#xff1a;为什么在Mac M系列芯片上跑Qwen-Image-Lightning是个“硬骨头”&#xff1f;Qwen-Image-Lightning&#xff0c;这个名字听起来像是一道闪电劈开图像理解的黑箱——它确实是通义千问团队推出的轻量级多模态模型&#xff0c;主打“快、小、准”&#xff1…

作者头像 李华
网站建设 2026/9/26 7:28:39

DeepSeek 接入 AI Agent 完全指南:新手快速上手的准备清单

DeepSeek 接入 AI Agent 完全指南&#xff1a;新手快速上手的准备清单 【免费下载链接】awesome-deepseek-agent 项目地址: https://gitcode.com/GitHub_Trending/aw/awesome-deepseek-agent awesome-deepseek-agent 是一份覆盖 Cherry Studio、Cline、Qwen Code、GitH…

作者头像 李华
网站建设 2026/9/26 7:28:24

AgentScope 2.0实战:多智能体协作框架的企业级落地指南

做多智能体应用开发半年多&#xff0c;我一直在找一套能让Agent们好好协作的框架。试过LangChain、AutoGen&#xff0c;也自己用消息队列拼过几套方案&#xff0c;总感觉差一口气——要么编排能力太弱&#xff0c;要么只适合Demo不适合生产。直到上个月把AgentScope 2.0完整跑通…

作者头像 李华
网站建设 2026/9/26 7:28:24

基于机器学习的音乐推荐系统:从日志解析到排序模型实战

简介&#xff1a;一套基于机器学习技术的音乐推荐系统完整项目&#xff0c;包含源代码与配套文档说明&#xff0c;面向计算机相关专业毕业设计、课程设计及期末大作业需求者&#xff0c;也适合希望进行项目实战的初学者。项目曾获导师认可并被评为高分项目&#xff08;评审99分…

作者头像 李华