☰
Flink实时推荐系统全链路拆解:从Kafka行为接入到Redis结果落地
2026/10/3 10:04:10 网站建设 项目流程

简介:一套基于Apache Flink的商品实时推荐系统完整项目源码,面向大数据方向的学生、开发者和推荐系统初学者。项目以Scala为主要语言,完整实现了从数据采集、预处理、特征工程到推荐算法、结果输出、系统优化的核心链路,并特别提供了Flink与HBase集成读写示例,以及Kafka模拟用户行为数据的生成脚本,便于在本地快速搭建可运行的实时推荐Demo。压缩包共44个文件,包含34个Scala源文件、SQL和HBase建表语句、Kafka数据模拟脚本、Maven工程配置等,整体体积仅245KB,内容紧凑却覆盖了完整工程结构。资源内部包含Flink读写HBase的模块,可观察向HBase写入数据的完整过程;SQL脚本提供用户表、商品表及行为表的建表语句,配合Kafka模拟数据脚本能快速生成测试数据流。已有275人学习下载,通过阅读源码和运行实践,可以深入理解Flink DataStream API、窗口统计、状态管理以及协同过滤等推荐算法在实时场景中的落地方法,是一份适合课程设计、毕业设计或项目实训的高质量参考资料。

1. 基于Flink商品实时推荐系统.zip:解压之后先别急着跑

拿到这份“基于Flink商品实时推荐系统.zip”,多数人的第一反应是解压、导入IDE、等Maven把依赖拉完,然后直接点运行。我的建议恰好相反:先别急着跑,跑不起来的。这个压缩包里不是一套能开箱即用的服务,而是一套“实时推荐链路的最小完整实现”——从Kafka里的用户行为日志,到Flink实时计算用户特征和商品特征,再到召回排序、写入Redis供API层读取。你需要的不是解压后的一瞬间,而是先把链路图画清楚,知道哪个环节缺了MySQL、Redis、Kafka或Druid,缺了之后会报哪类错误。这篇笔记就按“链路设计 → 落地步骤 → 参数调整 → 排错”往下讲,适合两类人:一类是拿这个zip做毕设或课程设计、需要把它改造成自己业务形态的在校生;另一类是刚接手实时推荐任务、想快速理解一套可用数据流的一线开发。无论哪种,先管住双击运行的手,我们从上到下把这条链路捋一遍。

2. 实时推荐链路设计:从行为埋点到召回排序的完整闭环

2.1 用户行为数据从哪来:Kafka + 日志格式的约定

Flink商品实时推荐系统处理的数据源头,通常是前端或客户端上报的曝光、点击、加购、下单四类行为报文。这些报文会统一打进Kafka的topic,常见的字段结构如下:

{ "userId": "u_10001", "itemId": "i_2048", "behavior": "click", "scene": "homepage_recommend", "timestamp": 1691234567890, "extra": { "duration": 3200, "page": "detail" } }

这个JSON结构里有两个字段对后续实时特征计算很重要。behavior是行为类型,推荐模型只应该学习正向行为,但曝光数据要单独保留,用于计算曝光过滤,避免把用户已经看过且没点击的商品反复召回。timestamp必须带,而且建议是毫秒级、事件时间语义的时间戳——后续做窗口聚合和浏览时长过滤时,如果用的是数据到达Flink的processing time,那么数据在kafka里积压一段时间后再消费,行为顺序会失真。

常见做法是把这四类行为分别路由到Kafka的不同分区甚至不同topic,实时Joiner再统一做流合并。这个zip里的工程多数只有click一种行为,我一般会在二次改造时把加购和下单单独拉出来,因为在推荐排序阶段,加购行为的权重应当远高于点击。

2.2 冷热链路分开:离线协同过滤与实时行为特征的接驳

商品推荐系统不可能只靠实时数据。实时数据解决的是“用户刚刚看了什么、我马上给他关联推荐”的短时兴趣,而离线协同过滤解决的是“和这位用户历史画像相似的人群喜欢什么”的长时兴趣。两者需要接驳,接驳点是Flink任务启动时的一次全量加载。

具体落地时,离线部分用Spark或Hive跑一个物品协同过滤,产出“相似商品对”表,比如item_id, similar_item_id, score,结果写入MySQL或可以直接读的HBase。实时任务启动时,先用RichCoFlatMapFunction把这些相似关系加载到内存中,形成一个Map<String, List<SimilarItem>>的映射。之后每来一条用户实时行为,就用当前行为商品ID去这个内存Map里捞相似商品,捞到的集合作为“协同过滤召回池”,再与实时热度池合并。这里有一个必须留意的边界:这个内存Map不能太大。几百万对相似关系大约占几百MB堆内存,可以接受;上亿对就会频繁Full GC,此时要换成外部缓存,把相似关系放到Redis里,实时任务改用异步IO查询。

提示:推荐系统里离线链路和实时链路的接驳,最忌讳在Flink任务里实时调MySQL查相似商品。每一条用户行为都打一次MySQL,行为量一上来,连接池先被打爆,然后Flink反压,最后Kafka消费延迟以分钟计。相似关系这种静态数据,启动加载或异步IO是两种可靠姿势。

2.3 实时特征的计算与写入:Redis里到底存什么

整套链路计算出的推荐结果,最终要落到API层能快速读取的存储里。Redis是这个zip项目最常见的落地存储,因为推荐结果的读取延迟要求通常在10毫秒内。写入Redis的数据结构是有讲究的,常见的两种设计如下。

第一种是按用户维度存一个推荐列表,使用Redis的String结构,key为rec:user:{userId},value为JSON数组,例如:

[ {"itemId": "i_3021", "score": 0.92, "reason": "clicked_similar"}, {"itemId": "i_1077", "score": 0.87, "reason": "hot"} ]

第二种是存一个带过期时间的ZSet,ZADD rec:user:{userId} 0.92 i_3021。两种结构的取舍取决于API层的读法。如果一次性返回十个商品,String结构更省事;如果需要在Redis端做分页或截断,ZSet更灵活。实际工程里我一般只写String结构,因为读取逻辑简单明了,且可以整体设置过期时间,比如EXPIRE rec:user:{userId} 1800——推荐结果只保证半小时内有效,用户再刷新一次页面时由新的实时计算重新生成。

这里的核心在于:实时推荐系统写入Redis的不是“模型训练出来的排序结果”,而是“当前时刻结合实时行为重新算过一轮的召回排序结果”。写Redis前要把itemId去重、过滤已曝光商品、按分数降序排列,这三个动作在Flink的Sink函数里完成比较合适。

3. 把zip里的工程跑起来:环境准备、自定义Source与Sink的落地步骤

3.1 从zip解压到Flink环境搭建:一次不用跑通全部的最小启动

先把环境搭好。这份工程依赖的组件最少有三个:Kafka、Redis、MySQL。MySQL用来存商品元数据与用户画像表,Kafka用来接收模拟行为数据,Redis用来存最终推荐结果。

# 启动一个本地Kafka(用docker compose是最省心的方式) # docker-compose-kafka.yml version: '3' services: zookeeper: image: confluentinc/cp-zookeeper:7.4.0 environment: ZOOKEEPER_CLIENT_PORT: 2181 kafka: image: confluentinc/cp-kafka:7.4.0 ports: - "9092:9092" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
docker compose -f docker-compose-kafka.yml up -d

Kafka起来之后,创建一个topic用于接收行为数据,然后确认可以生产消费:

docker exec -it kafka kafka-topics --create \ --bootstrap-server localhost:9092 \ --topic user_behavior \ --partitions 3 --replication-factor 1

partition数建议设成3。实时推荐任务在数据量不大时,单并行度也能跑,但3个分区可以让你在测试并行度调整时看到效果差异,也不用因为分区太多造成Flink checkpoint过大。如果zip里自带了模拟数据发送脚本,就先确认脚本往哪个topic发数据、发的数据格式是否和上文的JSON结构一致。这一步经常有偏差,脚本发的字段名和Flink解析的字段名对不上,运行时不会报错,但所有字段都是null,推荐结果永远是空列表。我一般会在这一步先用kafka-console-consumer消费几秒看看原始报文长相。

3.2 自定义DataSource:读取Kafka和模拟数据发生器的选择

这个zip的核心必然包含一个读取Kafka的自定义DataSource。Flink官方提供的FlinkKafkaConsumer已经能满足需求,但很多教学版zip喜欢写一个自定义的SourceFunction来模拟数据流,原因是可以脱离Kafka独立演示。这里的风险在于:演示用的自定义Source多半是无限循环生成随机JSON,而不是消费Kafka里的真实数据。跑通演示容易,上了生产环境还是得切回Kafka Connector。

// 自定义数据源:从Kafka读取用户行为JSON // DataSource与DataSink的自定义是理解Flink实时计算两条数据边界的核心 DataStream<String> rawStream = env.addSource( new FlinkKafkaConsumer<String>( "user_behavior", new SimpleStringSchema(), kafkaProps ) ); // 把JSON解析成JavaBean,这里用Flink自带的Jackson实现 DataStream<UserBehavior> behaviorStream = rawStream .map(new JsonToBehaviorFunction()) .returns(TypeInformation.of(UserBehavior.class)); // 关键参数:设置事件时间与水位线,保证后续窗口聚合的准确性 behaviorStream.assignTimestampsAndWatermarks( WatermarkStrategy .<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.getTimestamp()) );

这段代码里最值得调整的是forBoundedOutOfOrderness(Duration.ofSeconds(5))这一行。5秒是“允许乱序的最大时间”,如果Kafka里的行为数据时间戳偶尔乱序,5秒足够覆盖绝大多数场景。调大了会导致窗口结果晚出来,调小了会丢弃一部分乱序数据。流式推荐场景中,晚几秒出结果通常可以接受,丢数据反而会造成推荐结果空洞,所以这个参数我习惯于设到10秒。

3.3 自定义DataSink:结果写入Redis的可靠姿势

工程里自定义Sink是另一个必看模块。写入Redis的Sink需要自己实现连接池管理,不能每条结果都新建连接。连接池参数直接决定高吞吐下的稳定性,下面是最小可用的实现骨架:

public class RedisSink extends RichSinkFunction<RecommendResult> { private JedisPool jedisPool; @Override public void open(Configuration parameters) throws Exception { // 在open里初始化连接池,而不是在构造函数或每条数据里创建 JedisPoolConfig config = new JedisPoolConfig(); config.setMaxTotal(10); config.setMaxIdle(5); config.setMinIdle(2); config.setTestOnBorrow(true); // 连接池要复用,否则背压一来,连接先耗死 this.jedisPool = new JedisPool(config, "localhost", 6379, 3000); } @Override public void invoke(RecommendResult value, Context context) throws Exception { try (Jedis jedis = jedisPool.getResource()) { // 每个用户的推荐列表只保留一份,用JSON序列化后写入 String key = "rec:user:" + value.getUserId(); String json = objectMapper.writeValueAsString(value.getItems()); jedis.setex(key, 1800, json); } catch (Exception e) { // 写Redis失败不要直接抛出,先检查Redis是否可用 // 这里打印日志并跳过,比把整个任务搞挂更符合推荐场景 LOG.warn("write redis failed, userId={}", value.getUserId(), e); } } }

Sink里的两个细节要注意。一是setex同时设置了过期时间为1800秒,这比set更安全——真正生产环境里,如果Flink任务挂了,Redis里残留的旧推荐结果至少会自动过期,不会让用户一直看到几小时前的推荐。二是testOnBorrow(true)保证每次从池里借连接时验证连接是否存活。Redis服务正常时这个参数多一点点开销,但Redis重启过之后,这个参数可以避免大批连接报错。

提示:别在Sink里把异常直接抛出去。推荐结果的Sink属于“尽力而为”的写入,Redis短暂抖动导致几条写入失败,用户下一次请求时重算一次就可以了。相比之下,Kafka的offset提交和checkpoint稳定性更重要,这两个一旦失败会造成数据重复消费或丢失。

3.4 窗口聚合:热门商品榜为什么能反映实时热度

整个工程里最能体现“实时”二字的,是热门商品榜的计算。本质上是一个滑动窗口的计数聚合,Flink里用一行代码就能表达:

DataStream<ItemHot> hotStream = behaviorStream .filter(behavior -> behavior.getBehavior().equals("click")) .keyBy(UserBehavior::getItemId) .window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1))) .aggregate(new CountAggregate(), new HotWindowResult())

这段代码里窗口的两个时间参数决定了实时热度榜的敏感度。Time.minutes(10)是窗口长度,Time.minutes(1)是滑动步长,含义是每1分钟计算一次过去10分钟内的点击热榜。窗口越长,热度越平滑,短时冲高的商品不容易立刻上榜;窗口越短,热度越敏感,但容易出现某商品因一次小规模刷量就冲上榜的情况。商品推荐场景我一般用10分钟窗口、1分钟滑动,既保证对突发热点的响应及时,又不会让榜单一惊一乍。

窗口聚合的结果拿到之后,要和协同过滤召回、实时行为召回一起进入排序阶段。排序逻辑在推荐系统里可以很复杂,但在这份zip工程里通常就是一道加权公式的ProcessFunction。把点击、加购、下单分别给不同权重,再叠加热度分数,最后按总分排序取TopN。这一步虽然写法简短,但它决定了整个系统推荐结果的“是否像人推荐的”,比任何单独模块都值得反复调参。

4. 实时任务避坑指南:JDBC连接器异常、Kafka积压与数据落地的典型问题

4.1 JDBC连接器异常:周期性写失败不是网络问题,是连接池参数在裸奔

把实时计算结果同时写入MySQL用于离线分析,是这个zip里常见的扩展做法。但Flink JDBC连接器的报错频率,在各大数据平台的热搜词里居高不下,典型错误长这样:

Caused by: java.sql.SQLTransientConnectionException: HikariPool-1 - Connection is not available, request timed out after 30000ms

这个报错的现象是:任务运行前半小时正常,之后周期性出现写失败,而且每次失败间隔和窗口触发周期高度重合。原因基本是连接池的最大连接数配置小于并发写入的Task数量,窗口一批数据到来时,多个并行子任务同时申请连接,池里的连接被借完,新请求排队直到超时。

解决方法是调整JDBC连接器的参数:

-- 在Flink SQL中使用JDBC连接器时,通过SQL Hints调整连接池参数 INSERT INTO mysql_analytics_table /*+ OPTIONS('sink.buffer-flush.max-rows' = '200', 'sink.buffer-flush.interval' = '5s', 'sink.max-retries' = '3') */ SELECT * FROM realtime_features;

这里的三个参数各解决一个问题。sink.buffer-flush.max-rows控制攒够多少行才刷一次,默认是1000,但MySQL服务端如果性能一般,攒太多行一次写入容易造成锁等待;调低到200可以缓解。sink.buffer-flush.interval控制最多等多久必须刷一次,设5秒保证延迟可控。sink.max-retries是写入失败后的重试次数,默认值偏小,数据库抖动时容易直接把checkpoint搞失败。

注意:如果改完参数仍然周期性报错,去查MySQL的max_connections和wait_timeout。很多所谓“JDBC连接器异常”,根源是数据库侧的wait_timeout默认8小时,连接池里的连接长时间空闲被MySQL服务端断开,连接池却还认为连接存活。这属于经典的“两端参数不匹配”,排查方向从一开始就要把范围扩大到数据库服务端。

4.2 Kafka消息积压:反压源头往往不在Source,而在下游Sink

实时推荐任务有一个高频现象:Kafka消费延迟从秒级涨到分钟级,kafka-consumer-groups查看lag持续上涨。新手第一反应是加大并行度,于是把Source的并行度从3改到12,但很快发现lag不降反升。真实原因通常是下游有慢操作:要么是排序阶段里有外部调用,要么是Sink写入Redis时单条同步写太慢。

用Flink的Web UI排查时,看每个算子的BackPressure指标。如果Sink: RedisSink显示High,Source反倒正常,说明背压是从尾部往回传导的。此时解决方案不是加Source并行度,而是给Sink加批量写入能力:

// 解决方案:把逐条写Redis改成先攒一批再批量写 public class BatchRedisSink extends RichSinkFunction<RecommendResult> { private transient List<RecommendResult> buffer; private static final int BATCH_SIZE = 100; private static final long BATCH_INTERVAL_MS = 2000; @Override public void open(Configuration parameters) { this.buffer = new ArrayList<>(); } @Override public void invoke(RecommendResult value, Context context) { buffer.add(value); if (buffer.size() >= BATCH_SIZE) { flush(); } } private void flush() { // 批量写Redis,用pipeline减少RTT try (Jedis jedis = jedisPool.getResource()) { Pipeline pipeline = jedis.pipelined(); for (RecommendResult result : buffer) { pipeline.setex("rec:user:" + result.getUserId(), 1800, JsonUtils.toJson(result.getItems())); } pipeline.sync(); } buffer.clear(); } }

这个改动的关键点是pipeline。原本100条结果要100次RTT,pipeline可以合并成一次网络往返,吞吐量能提升一个数量级。代价是延迟从单条立即写入变成最多2秒的攒批窗口,对推荐结果而言完全可接受。做这类优化时,心里要有一杆秤:实时推荐系统的实时性,指的是行为发生后几秒内能反映到结果里,而不是每一条计算结果都即刻可查。

4.3 Sink Hive表数据不入表:分区提交是个假象

这个zip如果扩展了Hive数仓链路,会遇到一个很典型的“数据不入表”现象:flink任务在Web UI上显示sink成功,checkpoint也正常,但到Hive分区目录里看,.staging文件一堆,正式分区数据却是空的。这是把实时数据写入Hive表时的经典坑,原因基本可以锁定在“流式写入分区文件需要streaming阶段自动提交”。

排查顺序如下。先看表属性建表时是否声明了streaming相关参数。Hive表需要开启文件提交机制,否则Flink的FileSink会一直写.staging临时文件,永远不会rename成正式文件:

-- 建表时必须指定streming相关属性,缺了这个文件不会从staging转正 CREATE TABLE recommendation_log ( user_id STRING, item_id STRING, recommend_type STRING, ts BIGINT ) PARTITIONED BY (dt STRING, hh STRING) STORED AS PARQUET TBLPROPERTIES ( 'streaming' = 'true', 'auto-compaction' = 'false' );

再看Flink SQL或者DataStream写Hive时,是否设置了分区提交的触发策略。常见做法是:

// DataStream API写Hive时,开启分区提交的两种触发方式 // 1. 基于处理时间:每10分钟提交一次 // 2. 基于checkpoint:每个checkpoint尝试提交 FileSink<String> sink = FileSink.forBulkFormat( path, new ParquetRowDataBuilder(...) ) .withPartitionCommitter(new HivePartitionCommitter(conf, catalogTable)) .withPartitionCommitterStrategy(new MetastoreCommitPolicy()) .build();

很多教学工程只实现了写入逻辑,没有实现PartitionCommitter。这种情况下,数据确实写进了Hive的临时目录,Web UI也显示写出去了,但外表看不到任何数据。解决方案就是在FileSink上补上分区提交策略,同时把checkpoint间隔设置为与分区提交周期匹配的时长。如果不想动代码,最土的办法是定期用MSCK REPAIR TABLE修复分区,但不推荐上生产,治标不治本。

5. 用CEP实现“几分钟内浏览过A又浏览过B”的关联推荐:一个值得深度改造的进阶方向

5.1 为什么CEP比窗口聚合更适合商品关联推荐

协同过滤做的是“看了A的人还看了B”,实时热度做的是“现在大家都在看什么”,但这两者都缺一种能力:实时捕获“这位用户刚刚在短时间内连续浏览了哪些商品,把这些商品关联起来”。比如用户两分钟内依次查看了相机、镜头、三脚架,这时候最合理的推荐是把这三者的配件推荐出来。用滑动窗口聚合来处理这个场景非常别扭,因为窗口边界是死的;而CEP可以定义“A出现后10分钟内出现B”这类事件模式。Flink CEP的代码能直观表达这种业务规则。

5.2 一条可运行的CEP关联规则

// 定义模式:同一用户在10分钟内浏览了商品A,又浏览了商品B Pattern<UserBehavior, UserBehavior> pattern = Pattern .<UserBehavior>begin("first") .where(new SimpleCondition<UserBehavior>() { @Override public boolean filter(UserBehavior behavior) { return "click".equals(behavior.getBehavior()); } }) .next("second") .where(new SimpleCondition<UserBehavior>() { @Override public boolean filter(UserBehavior behavior) { return "click".equals(behavior.getBehavior()); } }) .within(Time.minutes(10)); // 在click流上应用这个模式,输出一次关联事件 DataStream<String> matchedStream = CEP.pattern(behaviorStream, pattern) .inProcessingTime() .select((Map<String, List<UserBehavior>> patternMap) -> { UserBehavior first = patternMap.get("first").get(0); UserBehavior second = patternMap.get("second").get(0); return first.getItemId() + "->" + second.getItemId(); });

这段代码里的within(Time.minutes(10))是“关联间隔”的核心参数。设太大,会把用户半小时前看过的商品和现在浏览的商品强行关联,关联噪音高;设太小,又捕捉不到真实的连续性浏览行为。我通常在电商场景先用5到10分钟起步,然后观察推荐结果里的关联商品点击率来反向调整。inProcessingTime()代表用处理时间做模式匹配,实时性高,但结果不可精确重放;如果后续要做效果对比实验,换成inEventTime()配水位线更严谨。

这块改造成本不高,但收益很明显:推荐结果里会出现一批“基于用户当前浏览序列”的动态关联商品。协同过滤提供的相似商品是相对静态的,而CEP关联规则让每个用户看到的是“跟着本次浏览路径走的推荐”,这是实时推荐最容易体现差异化的能力。

玩到这一步,这个zip工程已经不只是“能跑起来”,而是一个可以往生产形态演变的骨架了。回头看我自己的实践经历,每次做实时推荐改造,最深的教训都是同一个:不要把实时推荐当成一个纯计算问题,它是一个工程链路问题,Kafka里数据的质量、Redis里结果的过期策略、MySQL连接池的脾气,任何一个环节拉胯,Flink计算再正确也白搭。看不懂的线上故障,十有八九出在链路两端,而不是计算引擎本身。希望帮到你。

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询