ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

基于Flink的商品实时推荐系统:从Kafka到HBase的完整工程实践

基于Flink的商品实时推荐系统:从Kafka到HBase的完整工程实践 简介基于Flink的商品实时推荐系统完整代码包聚焦大数据实时计算与个性化推荐场景适合计算机专业学生、大数据开发初学者及对实时推荐感兴趣的从业者学习参考。压缩包内含四十四份文件以Scala源码为主共三十四份覆盖Flink作业核心逻辑另有SQL建表与查询语句、属性配置文件、HBase建表脚本、Kafka模拟数据生成脚本以及XML与TXT辅助说明可支撑从环境准备到数据接入、数据清洗、特征工程、推荐计算与结果输出的完整链路。包体约二百四十五KB目录层次清晰便于按功能模块阅读。目前已有二百七十五人学习下载。通过研读该项目可系统理解Flink窗口聚合、状态管理、检查点与保存点等机制在推荐系统中的应用掌握基于用户行为的特征提取方法以及协同过滤、矩阵分解等推荐算法的落地方式同时还能学习Kafka消息队列与HBase列式存储的集成用法是一份兼顾原理与实战的优质参考资源也适合作为课程设计与毕业设计的选题蓝本。1. 基于Flink的商品实时推荐点击后几百毫秒内换一批推荐这份工程源码全流程可跑用户前脚刚点开一个商品详情页后脚推荐位就换成了另一批关联商品中间通常只有几百毫秒。这套体验背后就是实时推荐系统在干活。我拆的这份《基于Flink的商品实时推荐系统》不是教学截图也不是伪代码zip里是完整的flink-recommend-system-main主工程、flink-2-hbase独立落地模块加上pom.xml、src源码、data示例数据和sql建表脚本一条链路从Kafka接日志、Flink做窗口统计、ALS出候选集、最后写进HBase全部能跑。适合刚把Flink官网菜鸟教程刷完、准备做课程设计或面试项目的人也适合离线推荐想迁实时但没有完整参考的开发。后面几章我把选型、启动、核心算子、坑点挨个拆开讲重点说清楚参数为什么那样设。2. 推荐系统为什么选Flink流计算选型逻辑与项目模块拆解2.1 用户行为不是宽表而是事件流实时推荐的引擎选型逻辑离线推荐是每天凌晨跑一次批任务把用户行为从MySQL或Hive捞出来算完协同过滤再写回推荐表。这套流程数据量不大时没问题但它天然是T1的用户今天中午点击的商品要等第二天凌晨才影响推荐结果。放到电商场景里用户看完一个商品后最可能继续买关联商品这个决策窗口只有几分钟等不到第二天。所以实时推荐必须有一个持续消费事件流、维护状态、低延迟输出的计算引擎。Flink在这个场景的核心优势有三点。第一是低延迟加高吞吐事件从Kafka进来到经过算子处理再到Sink写出秒级内完成背压机制避免数据堆积拖死下游。第二是原生状态管理用户浏览频次、购买记录都存在Keyed State里不需要每次计算都查一次数据库。第三是精确一次语义配合Checkpoint把状态和外部系统写入对齐推荐结果不会因为节点宕机重复或丢失。面试里常问“实时推荐为什么不用Spark Streaming”我的回答很直接Spark Streaming核心是按微批切数据本质还是批处理延迟在秒级到分钟级状态管理不如Flink原生。Flink把“流”当一等公民窗口、水位线、事件时间这些机制是搜“Flink实时计算进阶篇”时最常见的主题。这个工程选Flink不是跟风因为滑动窗口和状态管理正好命中推荐场景两个硬需求按时间窗统计用户短期偏好、跨窗口记住长期画像。2.2 六个环节一条链从日志采集到结果落库的数据流转整个实时推荐系统在架构上分成六个环节摘要里列得很清楚我按工程落地顺序重新梳理了一遍环节输入处理方式输出数据采集用户点击/浏览/购买行为日志Flink Source接入Kafka原始事件流数据预处理原始事件流清洗、去重、格式化标准行为流特征工程标准行为流滑动窗口统计频次、类别偏好用户特征向量推荐算法用户特征商品特征规则匹配/ALS评分候选商品集合结果输出候选商品集合HBase Sink持久化推荐ID列表系统优化全部运行期数据Checkpoint/状态TTL/并行度调优稳定运行的作业数据采集部分工程里用Kafka做消息中间件行为日志以JSON格式写入topicFlink的KafkaSource持续拉取。数据预处理是最容易被新手跳过的一步生产日志里几乎一定会混入爬虫流量、空字段事件和错乱时间预处理做去重、过滤异常IP、补全缺失字段。特征工程是实时推荐和离线推荐差异最大的地方离线是跑一条SQL把历史行为聚合好实时是Flink用滑动窗口滚动计算最近N分钟点击、最近N天购买。推荐算法环节分两步矩阵分解这类计算量大的在离线阶段完成Flink只加载模型参数做实时打分。结果输出单独拆了flink-2-hbase模块说明作者对落地这一层是认真处理的不是print到控制台就算完。六步连起来数据流长这样Kafka → FlinkSource → Filter/Map → KeyByWindow → 评分算子 → HBase Sink。我拆过不少类似工程这是实时推荐的标准骨架不管业务多复杂骨架基本不变。2.3 工程目录解读flink-recommend-system-main与flink-2-hbasezip解压后最外层是flink-recommend-system-main本身是Maven多模块工程。pom.xml在主目录下统管依赖版本flink-2-hbase是子模块单独负责推荐结果向HBase写入。src下是主代码data是示例数据sql是初始化脚本。目录/文件作用说明pom.xmlMaven父工程统管Flink、Kafka、HBase依赖版本flink-2-hbase落地子模块自定义HBase Sink、RowKey设计src/main/java主代码DataSource、预处理、窗口计算、推荐打分data示例数据用户行为日志、商品数据样例sql初始化脚本建表语句、基础数据插入我一般拿到工程先不看代码先用Maven把依赖树打出来确认版本兼容再打开sql目录看表结构。这个方法建议照做。pom.xml里的Flink版本和HBase客户端版本容易互相打架HBase 1.x客户端连HBase 2.x集群会出现RPC协议不兼容这类问题在第5章单独讲。目录层面有个值得注意的细节flink-2-hbase单独成一个模块意味着写HBase的逻辑和主计算逻辑解耦以后替换成写Redis或MySQL时不需要动主链路代码。接手别人实时推荐项目时这种模块边界就是代码里最值钱的架构信息。3. 从零把工程跑起来环境准备、SQL初始化和作业提交3.1 环境版本对照JDK、Maven、Flink、Kafka、HBase怎么配实时推荐系统涉及组件多版本错一个就可能出现ClassNotFoundException或连接超时。我先给一套实测下来兼容性比较稳的组合再说明每个组件在工程里承担的角色。组件推荐版本工程内作用JDK1.8编译运行主代码Maven3.6.x管理依赖与打包Flink1.13.x流式计算引擎Kafka2.8.x行为日志消息队列HBase1.4.x推荐结果存储MySQL5.7.x商品与用户维度数据安装顺序上建议先JDK再Maven这两个是前置工具。Flink解压就能用不需要安装守护进程无论local模式还是提交到YARN都通过bin/flink脚本操作。Kafka需要先启动Zookeeper再启动brokerHBase依赖HDFS或本地文件系统做存储MySQL放商品表和用户表推荐结果的主键关联也走它。这套组合不是唯一解Flink 1.13配老版本HBase客户端时要注意Netty依赖冲突细节在第5章展开。环境装好后用Flink自带的词频统计示例跑一遍验证安装这一步在官方叫“词频统计初体验”是最快确认环境没问题的办法。3.2 建表与基础数据sql目录里的脚本先跑哪张表sql目录下的脚本按“先维度表、后行为表、最后插入样例数据”的顺序执行。维度表包括用户表和商品表行为表记录点击、收藏、加购、购买事件。下面是核心三张表结构CREATE TABLE dim_user ( user_id BIGINT PRIMARY KEY, user_name VARCHAR(64), user_sex TINYINT, register_time TIMESTAMP ); CREATE TABLE dim_product ( product_id BIGINT PRIMARY KEY, product_name VARCHAR(128), category_id BIGINT, price DECIMAL(10,2), tags VARCHAR(256) ); CREATE TABLE user_behavior_log ( log_id BIGINT PRIMARY KEY, user_id BIGINT, product_id BIGINT, behavior_type VARCHAR(16), behavior_time TIMESTAMP, INDEX idx_user_time (user_id, behavior_time) );字段设计有几个关键点。behavior_type用字符串而不是数字因为“点击、收藏、加购、购买”在日志里本身是可读字符串排错时方便直接查真要压性能再考虑字典编码。behavior_time用TIMESTAMPFlink接入后能直接转成事件时间不丢毫秒精度。user_behavior_log上的联合索引user_id, behavior_time给离线查询和补数任务用实时链路不走MySQL但这个索引能加快排查数据问题时的手工查询。3.3 提交Flink作业local模式与YARN模式的命令参数工程编译打包后提交作业分两种常见形态。本地调试用local模式一条命令跑通全链路生产环境用YARN Application模式作业作为独立应用跑在集群上。下面是本地提交命令mvn clean package -DskipTests flink run -c com.recommend.streaming.UserBehaviorJob \ target/flink-recommend-system-1.0.jar \ --kafka.bootstrap.servers localhost:9092 \ --kafka.topic user_behavior \ --hbase.zookeeper.quorum localhost:2181 \ --hbase.table rec_user_product参数说明-c指定主类全限定名确保Flink找到作业入口。--kafka.bootstrap.servers是Kafka broker地址多个节点用逗号分隔。--kafka.topic是消费的原始行为topic。--hbase.zookeeper.quorum是HBase的Zookeeper地址客户端通过它拿到RegionServer列表。--hbase.table是推荐结果写入的目标表。这些参数在工程里一般通过ParameterTool解析对应主类里的一段固定代码。YARN模式下命令换成flink run -m yarn-cluster -yjm 1024 -ytm 2048-yjm是JobManager内存-ytm是TaskManager内存。实时推荐作业建议至少给TaskManager 2GB以上因为状态后端要存窗口聚合结果和用户画像。提交后打开Flink Web UI能看到各算子并行度、背压和水位线指标。作业起不来时能起作业但计算没输出问题在JobManagerTaskManager反复挂问题大概率在执行算子或者内存配置。4. 核心实现逐一拆解预处理、窗口统计、协同过滤与HBase数据落地4.1 预处理Filter和Map算子怎么清洗脏数据实时链路接到的数据基本是半结构化的爬虫脚本造的垃圾点击、字段缺失的行、时间戳为0的数据都得过滤。预处理用两个算子Filter过滤Map规范化。DataStreamUserBehavior cleanedStream rawStream .filter(new FilterFunctionString() { Override public boolean filter(String json) { if (json null || json.trim().isEmpty()) return false; if (!json.contains(\behavior_type\)) return false; if (json.contains(spider) || json.contains(robot)) return false; return true; } }) .map(new MapFunctionString, UserBehavior() { Override public UserBehavior map(String json) throws Exception { JSONObject obj JSON.parseObject(json); UserBehavior ub new UserBehavior(); ub.setUserId(obj.getLongValue(user_id)); ub.setProductId(obj.getLongValue(product_id)); ub.setBehaviorType(obj.getString(behavior_type)); ub.setBehaviorTime(obj.getLongValue(behavior_time)); return ub; } });逻辑说明filter里做三层判断空数据被过滤缺少行为类型的行被过滤带爬虫标识的行被过滤。map把JSON字符串转成Java对象同时完成类型转换behavior_time从毫秒时间戳变成Long类型方便后续直接作为事件时间使用。参数上两个地方要注意getLongValue在字段缺失时会抛异常我一般配合异常捕获把它当脏数据处理而不是让作业崩溃behavior_type不在预定义集合里时map阶段标记为UNKNOWN后面窗口计算单独处理不污染正常统计。4.2 滑动窗口行为统计最近N天购买频次的窗口参数特征工程里最常用的是滑动窗口统计。实时推荐要算“最近7天内每个用户对每个类目买过几次”这个统计如果手动维护状态会很麻烦Flink窗口直接解决。DataStreamUserBehavior keyedStream cleanedStream .assignTimestampsAndWatermarks( WatermarkStrategy.UserBehaviorforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, timestamp) - event.getBehaviorTime()) ); DataStreamTuple2Long, Long buyCounts keyedStream .filter(ub - purchase.equals(ub.getBehaviorType())) .keyBy(ub - ub.getUserId()) .window(SlidingEventTimeWindows.of(Time.days(7), Time.days(1))) .aggregate(new CountAggregate());逻辑说明assignTimestampsAndWatermarks把每条记录的behaviorTime提取为事件时间允许10秒乱序。10秒是按电商日志实际选的太大延迟高太小数据容易丢。keyBy按用户ID分key同一个用户的所有行为进入同一个子任务。SlidingEventTimeWindows.of(Time.days(7), Time.days(1))表示窗口长度7天、滑动步长1天含义是“每天算一次最近7天的购买统计”。aggregate里用自定义CountAggregate做增量聚合只保存一个计数状态不把窗口内所有数据缓存到内存这是实时作业不OOM的关键。这个窗口算出的频次作为用户偏好特征输出到下游评分算子。4.3 推荐算法落地离线训练ALS模型、在线实时评分的配合方式Flink生态里没有像Spark MLlib那样开箱即用的算法库这个工程推荐算法部分采用常见做法是“离线训练、在线使用”离线阶段用Spark MLlib或Python里的implicit库训练ALS矩阵分解模型得到用户因子矩阵和商品因子矩阵把因子向量导出到HDFS或MySQLFlink作业启动时把因子矩阵加载成广播变量每来一条用户行为事件实时计算用户向量与候选商品向量的点积评分。MapStateLong, float[] userFactorState ...; DataStreamScoredProduct scoredStream keyedStream .flatMap(new FlatMapFunctionUserBehavior, ScoredProduct() { Override public void flatMap(UserBehavior ub, CollectorScoredProduct out) { float[] userVector userFactorMap.get(ub.getUserId()); if (userVector null) return; for (Long productId : candidateProductIds) { float[] productVector productFactorMap.get(productId); float score dot(userVector, productVector); out.collect(new ScoredProduct(ub.getUserId(), productId, score)); } } });逻辑说明userFactorMap是广播进来的用户因子矩阵候选商品集合从维度表加载。flatMap里对每个用户行为事件遍历候选商品做向量点积。方案的精妙之处在于Flink只负责高频在线部分模型训练这种重计算完全绕开实时链路工程稳定性和准确率都可控。如果不想引入离线训练也可以退化成规则推荐用4.2节算出的购买频次把用户最近购买多的类目下的商品按价格区间排序输出效果差一些但代码量少很多。工程里两种方式都预留了接口拿UDF去替换评分逻辑即可。提醒一句别试图在Flink UDF里直接跑协同过滤训练每个并行子任务都做全量计算作业基本会被kill。4.4 自定义Data Sink写HBase的RowKey设计与批量参数flink-2-hbase模块的核心是自定义Sink直接把评分结果写入HBase。抱着“写进去就行”的心态会踩很多坑RowKey设计和批量写入参数是两个关键点。public class HBaseSink extends RichSinkFunctionScoredProduct { private Connection conn; private BufferedMutator mutator; Override public void open(Configuration parameters) throws Exception { org.apache.hadoop.conf.Configuration conf HBaseConfiguration.create(); conf.set(hbase.zookeeper.quorum, localhost:2181); conn ConnectionFactory.createConnection(conf); BufferedMutatorParams params new BufferedMutatorParams(TableName.valueOf(rec_user_product)) .writeBufferSize(6 * 1024 * 1024); mutator conn.getBufferedMutator(params); } Override public void invoke(ScoredProduct sp, Context ctx) throws Exception { byte[] rowKey Bytes.toBytes(sp.userId _ sp.productId); Put put new Put(rowKey); put.addColumn(Bytes.toBytes(rec), Bytes.toBytes(score), Bytes.toBytes(String.valueOf(sp.score))); mutator.mutate(put); } Override public void close() throws Exception { if (mutator ! null) mutator.flush(); if (conn ! null) conn.close(); } }逻辑说明open阶段初始化HBase连接writeBufferSize设置6MB的批量缓冲区积累到6MB或到达flush周期才批量写RegionServer。invoke阶段每来一条评分结果做一次Put交给BufferedMutator不立即写入。close阶段必须flush一次确保作业结束时缓冲区里剩余写请求落地。RowKey用user_id加product_id拼接查询“某个用户推荐了哪些商品”时按user_id前缀做Scan避免全表扫描。批量参数三个需要调writeBufferSize太小批量效果差太大会让单次写入内存压力过高写入线程数默认即可HBase的hbase.client.write.buffer配合mutator的BufferSize一起生效。这条Sink如果改成写Redis就是另一套映射但结构完全一致这也是我建议你读懂它的原因——它是推荐结果落地层的样板代码。5. 避坑线上翻车最多的五个实时推荐问题与排查路径实时推荐系统的坑大部分不在“算法不精准”而在“作业跑不稳”。下面5条从实际运行里沉淀下来每条按现象、原因、解决三步讲。5.1 窗口迟迟不触发事件时间与水位线的“玄学”现象作业提交成功Kafka正常消费控制台也有日志输出但窗口聚合结果一直不出来Flink Web UI里窗口触发次数为0。原因事件时间模式下窗口必须等水位线超过窗口结束时间才会触发。代码里没有assignTimestampsAndWatermarks或日志里的behavior_time是处理时间而非日志产生时刻水位线都不会推进。另一个常见原因是数据本身乱序但没有给forBoundedOutOfOrderness预留延迟数据到达时窗口已经关闭。解决给数据流显式分配事件时间设置合理乱序容忍度。先用Web UI的watermark列观察水位线是否推进不推进优先检查时间戳分配逻辑推进了仍不触发再排查时间字段是否真的按毫秒填的。这一步是排查入口比盲改窗口参数有效得多。5.2 JDBC连接器频繁报错驱动类位置与连接池参数现象作业运行一段时间后抛出SQLException报连接超时或Connection is not available重启后恢复过几小时又挂。部分情况在JdbcSink写MySQL这一步高频出现很多人搜“flink的jdbc连接器异常”时都会撞上同款。原因JdbcSink默认每条数据获取连接、写入、释放高吞吐下连接池被频繁打满。另一个隐蔽原因是MySQL驱动jar没打进作业包运行时抛ClassNotFoundException。这个在本地IDEA里跑不出来因为IDEA从Maven仓库能找到驱动打成的fat jar里却没有。解决驱动问题把mysql-connector-java以compile方式放进pom依赖打包后确认打成了可执行jar。连接池问题给JdbcExecutionOptions设置参数batchSize设1000条批量提交batchIntervalMs设15秒兜底同时把连接池最大连接数调成跟Flink并行度一致。参数改完观察Web UI的Sink端背压指标背压变黄说明写入瓶颈在连接池而不是下游MySQL。5.3 HBase数据半天不落地异步Buffer没flush的血泪经验现象作业状态正常HBase表也建了查RegionServer发现数据量不变。等很久后突然一批数据出现而且每次都是攒一大波。类似的问题在Sink到Hive时也常见表建好了数据不进表多半是流式写入的batch开关没配对。原因用BufferedMutator异步写入时writeBufferSize设得过大而数据量不足以撑满缓冲区mutator不会主动flush数据一直积在内存里。另一个翻车点是close时没调flush作业在流式场景长期不close数据自然永远不落库。解决把writeBufferSize调到业务可接受的量比如6MB同时给HBase客户端设置定时刷新参数。更稳的做法是缓冲到一定数量或一定时间就调用一次mutator.flush()用“条数加时间”双阈值控制数据延迟。从那以后我写HBase Sink都强制走一遍“close里必须flush、buffer里必须配定时刷新”的检查。5.4 Kafka重启后重复消费消费者组offset提交时机现象作业重启后推荐结果和重启前完全重复用户看到的推荐位内容“倒退”了。如果Flink作业配套了MySQL或Redis消费进度记录还会出现两边进度对不上。原因Flink的Kafka connector默认开启Checkpoint后offset由Checkpoint保存重启时从最近一次Checkpoint恢复。没开Checkpoint或者enable.auto.commit被误设成true且提交时机跟实际处理进度不同步重启后会按自动提交的offset恢复造成重复消费。解决开启Flink Checkpoint间隔60秒重启时从Checkpoint状态恢复offset以Checkpoint里保存的为准。排查看两个地方一是代码里有没有env.enableCheckpointing二是Web UI的Checkpoints页面有没有成功历史记录。确认开了Checkpoint仍重复消费再查Kafka topic副本写入参数和是否使用事务型生产者事务型生产者能配合Flink两阶段提交保证精确一次语义。5.5 状态无限膨胀OOM没有给状态设置TTL的后悔药现象作业运行一周左右TaskManager频繁内存溢出GC越来越频繁最终Container被杀。日志能看到OutOfMemoryError堆转储分析发现大量用户历史行为对象占满堆。原因Keyed State里存了用户维度的所有历史行为且没设置StateTtlConfig状态只增不减。真实流量下用户行为量无限增长几百万用户、每个用户几千条行为状态后端撑爆只是时间问题。这也是实时推荐区别于词频统计初体验的地方——词频作业结束就丢状态推荐系统是长驻作业状态生命周期必须显式管理。解决给StateDescriptor配TTL用户最近7天的行为统计设7天过期配合清理策略。StateTtlConfig ttlConfig StateTtlConfig .newBuilder(Time.days(7)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); ValueStateDescriptorLong stateDesc new ValueStateDescriptor(buyCount, Long.class); stateDesc.enableTimeToLive(ttlConfig);注意TTL生效需要时间而且是懒清理内存已经告急再配TTL已来不及。所以要在作业上线前把TTL配好TaskManager内存留出30%以上余量。排查内存问题时先看Web UI的State Size指标判断是状态增长还是算子数据堆积处理方案完全不同。6. 进阶用火焰图定位热点算子、把推荐链路血缘管起来6.1 火焰图定位算子热点作业跑顺之后怎么判断哪个算子占用CPU最多Flink不自带火焰图但可以用async-profiler对TaskManager进程采样把CPU消耗对应到算子和业务代码行上。常见做法是给TaskManager加JVM参数-agentpath:/opt/async-profiler/build/libasyncProfiler.sostart,eventcpu,fileflink_profile.html采样5到10分钟后下载生成的html火焰图在浏览器里看栈帧宽度。栈顶宽度最大的方法就是热点。推荐系统里最常见的火焰图形态是序列化和正则匹配占大头序列化热点把JSON解析换成Avro或更快的解析方式正则匹配热点把过滤里的“spider|robot”这类正则换成前缀匹配。这个方法好在能在“猜”和“测”之间给出第三方案用数据定位再改代码。6.2 把Flink推荐链路的血缘管起来实时推荐链路涉及Kafka、Flink、HBase、MySQL多个系统数据从哪来、经过哪些算子、产出到哪张表如果只靠人肉记三个月后没人说得清这是典型的“黑匣子”问题。用OpenMetadata这类元数据工具可以采集Flink作业的血缘关系把Kafka topic、算子、HBase表串成一条可视化链路。做实时数仓的团队常用这手应对资产盘点。它跟推荐系统结合的价值在于当某张HBase表推荐结果突然异常沿血缘能快速定位是Kafka字段变更还是Flink逻辑问题。这块我个人的习惯是新接手的Flink工程第一周先跑通第二周就补血缘元数据和关键指标监控不要等出问题再补。以前接过一个线上推荐作业团队没人知道某个字段是哪个上游加的排查花了两天最后靠翻git历史才理清。从那以后每次部署实时作业我都强制走一遍“代码评审、状态TTL检查、血缘注册”三重检查能拦下大多数线上事故。这份基于Flink的商品实时推荐系统zip按第3章的流程启动、把第4章的算子逐个改参数看效果两个小时内能跑出第一批推荐结果希望帮到你。本文还有配套的精品资源点击获取
返回列表