基于Flink与HBase的实时推荐系统实战:从源码到SpringBoot接口

📅 发布时间:2026/10/9 8:06:45
基于Flink与HBase的实时推荐系统实战:从源码到SpringBoot接口
简介这是一套基于Flink的商品实时推荐系统完整项目资料面向计算机、大数据、人工智能等相关专业的在校学生、教师及企业开发者可用于毕业设计、课程设计、项目立项演示或技术进阶学习。资源包共47个文件约247KB以34个Scala源码为核心配合SQL脚本、properties配置、XML依赖文件、HBase建表语句、Kafka模拟数据及说明文档覆盖实时推荐系统的数据采集、流式计算与存储落地等关键环节。项目已通过测试运行功能完整并获导师认可答辩评审达95分。读者可据此理解Flink实时计算与推荐逻辑的整合方式掌握HBase存储与Kafka数据模拟的配合思路也可在现有代码基础上修改扩展实现个性化功能。目前已有57人学习关注适合希望快速上手大数据实时推荐项目、积累工程实践经验的学习者参考使用。1. 从一份 Flink 商品实时推荐系统源码包说起它到底能跑出什么电商场景里最典型的实时链路是什么用户点一下商品后台要在几百毫秒内更新推荐列表。离线跑批那套 T1 的玩法早就撑不住这种需求了。这份基于Flink商品实时推荐系统详细文档全部资料.zip就是冲着这个场景来的——它把 Flink 实时计算、HBase 存储、推荐算法串成了一条完整链路附带详细文档和全部源码解压后能看到flink-recommend-system-main、flink-2-hbase、pom.xml、src、data、sql这些目录结构清晰不是那种扔一堆 class 文件让你自己猜的包。它适合谁正在做大数据课程设计、毕业设计的学生需要快速搭一个实时推荐 demo 的初中级工程师以及想理解 Flink 与 HBase 怎么配合做在线特征存储的开发者。不适合谁指望开箱即用直接上生产的人——这是教学级项目不是工业级推荐中台。但作为理解「实时推荐链路怎么搭」的参照物它的完整度足够让你少走很多弯路。下面我从环境搭建一路拆到调优把这份资源真正用起来。2. 环境搭建与工程结构拆解从 pom.xml 到数据流向2.1 先看清工程骨架再动手解压后的目录结构不是随便摆的每一层都有明确分工。flink-recommend-system-main是主工程目录里面包含 Flink 作业的入口和核心计算逻辑flink-2-hbase是 Flink 写入 HBase 的连接模块负责把计算结果落到 HBase 做在线查询data目录放的是模拟用户行为日志和商品元数据sql目录里是建表语句包括 HBase 表结构和可能用到的 MySQL 维表。pom.xml分两层——根目录的父 pom 管理版本子模块的 pom 声明各自依赖。我一般拿到这种包第一件事不是急着mvn compile而是先看父 pom 里的properties段确认 Flink 版本、HBase 版本、Scala 版本这三个关键坐标。版本对不上后面编译报错能让你查半天。常见做法是打开根目录pom.xml找到类似下面的配置properties flink.version1.14.4/flink.version hbase.version2.3.7/hbase.version scala.binary.version2.12/scala.binary.version java.version1.8/java.version /properties这三个版本号决定了你本地要装什么。Flink 1.14 对应 Scala 2.12HBase 2.3.x 需要 JDK 8。如果你本地是 JDK 11 或 17编译时可能遇到模块访问权限问题建议直接用 JDK 8 跑通再说。flink-2-hbase模块里通常会有 HBase 的Connection配置类里面写死了 ZooKeeper 地址和端口这个后面要改。2.2 本地环境准备与依赖拉取跑这个项目需要三样东西JDK 8、Maven 3.6、一个能用的 HBase 实例。Flink 本身不需要单独安装项目里用的是 Flink 的本地执行环境提交到本地 JVM 就能跑。但 HBase 必须真实存在因为代码里会去连它。HBase 的搭建方式有两种直接用 Docker 拉一个单机版或者在本地装一个伪分布式。我倾向于 Docker干净且不污染宿主机环境。命令如下docker run -d --name hbase-standalone \ -p 2181:2181 -p 16000:16000 -p 16010:16010 -p 16020:16020 -p 16030:16030 \ harisekhon/hbase:2.3端口映射说明2181 是 ZooKeeper16000 是 HBase Master 的 RPC 端口16010 是 Master 的 Web UI16020 和 16030 分别是 RegionServer 的 RPC 和 Web UI。跑起来后访问http://localhost:16010能看到 Master 界面就说明 OK 了。接下来在项目根目录执行依赖拉取mvn clean install -DskipTests-DskipTests是必须的因为项目里可能包含需要连接外部服务的测试用例本地环境没完全就绪时会卡住。如果拉取过程中卡在某个依赖上检查 Maven 的settings.xml是否配了国内镜像。常见做法是加阿里云镜像速度会快很多。2.3 数据流走向与核心模块职责这个项目的实时链路可以概括为数据源模拟用户行为 → Flink 消费并计算 → 结果写入 HBase → 对外提供查询。data目录里的日志文件通常包含userId、itemId、behaviorType、timestamp这几个字段模拟的是点击、收藏、加购这些行为。Flink 作业里一般会做这几件事先用SourceFunction或者KafkaSource读数据然后按用户或商品分组做窗口聚合算出热度分或者协同过滤的中间结果最后通过flink-2-hbase模块的 Sink 写入 HBase。sql目录里的建表语句对应 HBase 的namespace和table表名通常是user_behavior或item_similarity这类。提示如果你本地没有 Kafka项目里大概率提供了基于集合的SourceFunction作为替代在src里搜FromElements或fromCollection能找到。先用这个跑通逻辑再换成 Kafka 接真实数据。3. 核心链路跑通从行为日志到 HBase 推荐结果3.1 数据源接入与 Flink 作业启动项目里数据源的接入方式决定了你启动作业时要传什么参数。如果用的是SourceFunction模拟数据直接在 IDE 里右键运行主类就行如果用的是 Kafka需要先确保 Kafka 服务在跑并且 topic 名称和代码里一致。我一般会先找到主类通常在flink-recommend-system-main模块的src/main/java下类名可能叫RecommendJob或Main。看它的main方法里有没有StreamExecutionEnvironment.getExecutionEnvironment()有的话就是标准 Flink 作业入口。启动前检查data目录的路径是否写死在代码里如果是相对路径工作目录要设对。StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 本地调试设为1方便看输出 DataStreamString behaviorStream env.addSource(new BehaviorSource()); // BehaviorSource 从 data 目录逐行读取模拟日志setParallelism(1)在本地调试时很关键并行度大于 1 会导致输出交错你根本看不清哪条数据对应哪个结果。等逻辑验证通过后再调大并行度。3.2 推荐计算逻辑与 HBase 写入推荐计算部分是这个项目的核心。常见实现有两种一种是基于物品的协同过滤算 item 之间的相似度另一种是基于热度的简单推荐统计每个商品的点击量然后排序。教学项目里后者更常见因为逻辑简单、容易验证。计算完的结果会通过flink-2-hbase模块写入 HBase。这个模块里通常有一个HBaseSink类继承RichSinkFunction在invoke方法里调用 HBase 的Put操作。关键参数是表名和列族名这两个必须和sql目录里的建表语句对上。public class HBaseSink extends RichSinkFunctionRecommendResult { private Connection connection; private Table table; Override public void open(Configuration parameters) throws Exception { Configuration conf HBaseConfiguration.create(); conf.set(hbase.zookeeper.quorum, localhost); conf.set(hbase.zookeeper.property.clientPort, 2181); connection ConnectionFactory.createConnection(conf); table connection.getTable(TableName.valueOf(recommend_result)); } Override public void invoke(RecommendResult value, Context context) throws Exception { Put put new Put(Bytes.toBytes(value.getUserId())); put.addColumn(Bytes.toBytes(info), Bytes.toBytes(itemIds), Bytes.toBytes(String.join(,, value.getItemIds()))); table.put(put); } }open方法里初始化连接invoke方法里逐条写入。注意hbase.zookeeper.quorum要改成你实际的地址Docker 跑的话就是localhost。表名recommend_result需要提前在 HBase 里建好建表语句在sql目录里找。3.3 验证结果是否真正落库作业跑起来后别只看控制台没报错就以为成功了。要去 HBase 里查数据。用 HBase Shell 执行docker exec -it hbase-standalone hbase shell进入 Shell 后scan recommend_result, {LIMIT 5}能看到rowkey是用户 ID列族info下有itemIds列值是一串逗号分隔的商品 ID就说明链路通了。如果scan出来是空的先检查 Flink 作业的日志里有没有TableNotFoundException或Connection refused这两个是最常见的。注意HBase 的scan默认返回所有版本如果表里数据量大加LIMIT避免刷屏。另外rowkey设计如果用的是用户 ID注意热点问题——教学项目数据量小无所谓但心里要有这根弦。4. 避坑与排查那些让你卡半天的常见问题4.1 启动报 NoClassDefFoundError 或 ClassNotFoundException现象mvn clean install成功但 IDE 里运行主类时抛NoClassDefFoundError通常指向 HBase 或 Flink 的某个类。原因Maven 依赖作用域问题。Flink 和 HBase 的依赖在pom.xml里可能被标成了provided编译时能过运行时 JVM 找不到。教学项目里经常这么写因为默认你会用flink run提交集群环境有这些 jar。解决在 IDE 运行配置里勾选 Include dependencies with Provided scope或者临时把pom.xml里的scopeprovided/scope注释掉重新mvn install。跑通后再改回去。4.2 HBase 连接超时或 ZooKeeper 拒绝连接现象作业启动后卡在open方法日志里反复出现Connection refused或Session expired。原因HBase 没启动或者hbase.zookeeper.quorum配的地址不对。Docker 跑 HBase 时容器内的 ZooKeeper 监听的是容器 IP不是localhost但端口映射到宿主机后宿主机用localhost访问是通的。问题往往出在代码里写的是容器名而不是localhost。解决确认conf.set(hbase.zookeeper.quorum, localhost)和conf.set(hbase.zookeeper.property.clientPort, 2181)这两行。如果还不行进容器docker exec -it hbase-standalone bash执行hbase shell看能否本地连上排除 HBase 自身问题。4.3 写入 HBase 时报 TableNotFoundException现象Flink 作业日志里抛org.apache.hadoop.hbase.TableNotFoundException: recommend_result。原因表没建。sql目录里的建表语句没执行或者执行时 namespace 不对。解决进 HBase Shell先list看有哪些表。如果没有目标表找到sql目录里的.hbase或.txt文件把create recommend_result, info这行贴进 Shell 执行。注意表名和列族名要和代码里完全一致大小写敏感。4.4 数据写入成功但查询为空现象Flink 日志显示 Sink 正常执行没有异常但scan查不到数据。原因rowkey设计问题或者Put的列族写错了。比如建表时列族叫info代码里写成了dataHBase 不会报错但数据写到了不存在的列族下scan默认只显示已存在的列族。解决describe recommend_result看列族名和代码里的Bytes.toBytes(info)对比。另外检查rowkey是否为空——如果userId是 nullPut会抛异常但有些版本的 HBase 客户端会静默吞掉。4.5 本地跑通但换环境就挂现象在自己电脑上跑得好好的换到实验室服务器或云主机上就各种报错。原因路径写死、IP 写死、JDK 版本不一致。data目录的相对路径在不同工作目录下解析结果不同localhost在容器环境里指向容器自身而非宿主机。解决把data路径改成可配置的通过ParameterTool从启动参数传入。HBase 地址也抽成配置项。JDK 版本用java -version确认和pom.xml里的java.version对齐。5. 进阶玩法把推荐结果接到 SpringBoot 查询接口5.1 为什么要在 Flink 和前端之间加一层 SpringBootFlink 作业把推荐结果写进 HBase 后前端或者移动端不能直接连 HBase——协议不兼容权限也不好控。常见做法是在 HBase 前面加一个 SpringBoot 服务对外暴露 REST 接口内部用 HBase Client 查数据。这样 Flink 只管算和写SpringBoot 只管读和返回 JSON职责清晰。这个项目本身没包含 SpringBoot 模块但你可以基于flink-2-hbase里的 HBase 连接代码快速搭一个。我一般会新建一个 Maven 模块引入spring-boot-starter-web和 HBase Client把HBaseSink里的Connection初始化逻辑复用过来改成PostConstruct初始化。5.2 查询接口的实现与参数设计接口设计上至少提供一个按userId查推荐列表的 GET 接口。参数就一个userId返回 JSON 数组。代码如下RestController RequestMapping(/api/recommend) public class RecommendController { private Connection connection; PostConstruct public void init() throws IOException { Configuration conf HBaseConfiguration.create(); conf.set(hbase.zookeeper.quorum, localhost); conf.set(hbase.zookeeper.property.clientPort, 2181); connection ConnectionFactory.createConnection(conf); } GetMapping(/{userId}) public ListString getRecommend(PathVariable String userId) throws IOException { Table table connection.getTable(TableName.valueOf(recommend_result)); Get get new Get(Bytes.toBytes(userId)); Result result table.get(get); byte[] value result.getValue(Bytes.toBytes(info), Bytes.toBytes(itemIds)); if (value null) { return Collections.emptyList(); } return Arrays.asList(Bytes.toString(value).split(,)); } }PostConstruct保证Connection在服务启动时初始化一次避免每次请求都建连接。get操作按rowkey精确查询比scan快得多。返回的itemIds是逗号分隔的字符串拆成 List 后 SpringBoot 自动序列化成 JSON 数组。5.3 联调时容易忽略的两个细节第一个是 HBase 的Connection是线程安全的但Table不是。上面的代码每次请求都getTable在高并发下会有性能损耗。更好的做法是把Table也缓存起来或者用Connection.getTable的线程安全版本。教学 demo 无所谓但心里要知道这个边界。第二个是跨域问题。如果前端页面和 SpringBoot 服务不在同一个端口浏览器会拦请求。加一个CrossOrigin注解或者全局 CORS 配置就行。我一般直接在 Controller 上加CrossOrigin(origins *)本地调试省事。CrossOrigin(origins *) RestController RequestMapping(/api/recommend) public class RecommendController { // ... }联调顺序建议先确保 Flink 作业在跑并且 HBase 里有数据再启动 SpringBoot 服务最后用curl或 Postman 测接口。curl http://localhost:8080/api/recommend/1001能返回商品 ID 列表整条链路就算闭环了。从那以后我每次拿到这种实时推荐项目都强制先跑通「数据源 → Flink → HBase → 查询接口」这一条最小闭环再去调算法和参数。顺序反了后面全是玄学问题。希望帮到你。本文还有配套的精品资源点击获取