☰
Flink实时推荐系统:架构、核心算子与生产避坑实践
2026/10/3 9:04:37 网站建设 项目流程

简介:基于Flink的商品实时推荐系统完整项目,面向具备一定大数据基础、希望掌握实时推荐全流程的计算机学生与开发者。资源覆盖数据采集、预处理、特征工程、推荐算法到结果输出等核心环节,并附有HBase建表语句与Kafka模拟数据,便于快速跑通实验环境。压缩包共44个文件,以scala源码为主(34个),另有xml配置、sql脚本、properties配置与txt说明,整体245KB,代码规模精简、模块划分清晰,适合用于课程设计、毕业设计或Flink入门实战。已有275人学习下载。通过该系统源码,可直观理解Flink DataStream API、窗口统计与状态管理在用户行为处理中的应用,掌握协同过滤或矩阵分解生成实时推荐结果的方法,同时熟悉项目目录结构、自测运行路径与常见调参思路,是将大数据技术落地为推荐工程的实用参考。

1. 拿到“基于Flink商品实时推荐系统.zip”之后:先拆链路,别急着解压跑代码

一个基于 Flink 商品实时推荐系统的压缩包,打开之后通常是这副光景:README 写着环境要求,src 下面躺着 source、process、sink 三个包,data 目录放几份小样本行为日志。我的建议是,先别解压跑代码,先把链路拆出来:用户行为从哪进,特征算完放哪,召回结果给谁读。这类项目真正值钱的不是算分公式,而是从 Kafka 到 Flink 再到 Redis 这条实时管道能不能在秒级内把行为变成推荐依据。

这个方向适合两类人:一类是数据工程师,想在自己的平台上接一套实时特征计算;另一类是准备转实时计算方向的候选人,需要把窗口、状态、背压这些概念落到一个能讲的完整项目里。读完你会清楚每一层为什么这样选,以及任务延迟、数据不入库这类事故要从哪开始查。

2. 把 Flink 环境跑起来:从本地启动到用自定义 DataSource 喂行为数据

先解决“能不能跑”的问题。很多同学拿到压缩包第一件事就是把 pom.xml 里的依赖版本改一遍,结果依赖冲突直接把人劝退。我的建议是分两步:先让一个空作业跑通,再逐步把推荐逻辑加回去。下面按我平时搭环境的顺序来。

2.1 环境选型:本地 vs Standalone vs YARN,最小提交命令怎么写

推荐项目通常涉及三个运行环境:本地 IDE 调试、Standalone 集群验证、YARN 或 K8s 部署。我一般建议先本地 IDE 跑通,再用 Standalone 模拟生产,最后再谈资源调度。Flink 安装配置到部署这步,本地玩最容易翻车的地方是内存参数没调,窗口开大了直接被容器杀掉,看起来像代码问题,其实是资源配置问题。

先看本地模式。解压安装包后改动一个配置文件conf/flink-conf.yaml,最值得提前设的是这几项:

taskmanager.memory.process.size: 2048m taskmanager.numberOfTaskSlots: 4 parallelism.default: 4 state.backend: rocksdb state.checkpoints.dir: file:///data/flink/checkpoints

taskmanager.memory.process.size是整个 TaskManager 进程的总体内存,包括堆外和网络缓冲,只调堆内存不够。state.backend: rocksdb在本地看起来多余,但推荐项目一旦要做窗口加状态,RocksDB 能避免堆内存被状态撑爆,这个习惯最好一开始就养好。state.checkpoints.dir若是本地路径,重启机器后 checkpoint 会丢,后面第 4.4 节会讲到它的坑。

Standalone 模式下,提交任务的最小命令长这样:

bin/start-cluster.sh mvn clean package -DskipTests bin/flink run -d -p 4 \ -c com.recommend.job.RecommendationJob \ target/recommend-1.0.jar

-d表示分离模式,任务在后台跑,好处是你的终端不会一直挂着日志;坏处是任务异常退出时你不会第一时间收到堆栈。所以我本地调参数时反而不加-d,直接前台跑,让日志打到当前 session。-p 4是并行度,它覆盖配置文件里的parallelism.default,两者不一致时以命令行参数为准。-c指定主类,压缩包里如果重写了 main,或者多个类都有 main,这一步最容易踩坑。

如果你在 IDE 里直接跑,不需要启动集群。直接在env.execute()前写好StreamExecutionEnvironment.getExecutionEnvironment(),它会在本地自动创建一个 mini 集群。我见过不少人拿着flink run命令去 IDE 里找入口,绕了一大圈,其实 IDE 运行只需要 JDK 和 Maven 依赖。

2.2 自定义 DataSource:让模拟行为流变得可控

压缩包里常见的样子是自定义 DataSource,而不是直接读 Kafka。因为本地跑推荐项目,没有真实埋点数据,需要一个能控制速率和规模的模拟数据源。下面的代码是这类项目的典型写法:

public class BehaviorSource extends RichSourceFunction<BehaviorEvent> { private volatile boolean running = true; private long maxEvents; private long sleepMs; public BehaviorSource(long maxEvents, long sleepMs) { this.maxEvents = maxEvents; this.sleepMs = sleepMs; } @Override public void run(SourceContext<BehaviorEvent> ctx) throws Exception { long emitted = 0; Random random = new Random(); int[] userIds = {1001, 1002, 1003, 1004}; int[] itemIds = {5001, 5002, 5003, 5004, 5005}; String[] behaviors = {"view", "click", "cart", "buy"}; while (running && emitted < maxEvents) { BehaviorEvent event = new BehaviorEvent(); event.setUserId(userIds[random.nextInt(userIds.length)]); event.setItemId(itemIds[random.nextInt(itemIds.length)]); event.setBehavior(behaviors[random.nextInt(behaviors.length)]); event.setTimestamp(System.currentTimeMillis()); ctx.collect(event); emitted++; Thread.sleep(sleepMs); } } @Override public void cancel() { running = false; } }

注意这里我继承了RichSourceFunction而不是基础SourceFunction,原因是 Rich 版本能拿到RuntimeContext,后面你要在 source 里读配置、访问广播状态,甚至做个简单的限流,能力都在。ctx.collect(event)是真正的发射点,Flink 会把对象交给下游算子。Thread.sleep(sleepMs)控制生成速率,这一行非常关键,它决定了你的窗口计算能不能观察到滑动效果。

用的时候这样加进作业:

DataStream<BehaviorEvent> stream = env .addSource(new BehaviorSource(1_000_000L, 10L)) .setParallelism(2) .name("user-behavior-source");

setParallelism(2)和作业并行度不等同。source 并行度太高时,下游窗口算子容易接收乱序数据,处理时间窗口还好,事件时间窗口就要配合水印。name("user-behavior-source")不是摆设,它会在 Web UI 和火焰图里显示,第 4.3 节排查背压时你就能直接定位是哪个算子叫这个名字的资源卡住了。

2.3 行为数据入口:Kafka 流与 MySQL CDC 的取舍

本地跑通了,接下来要思考生产接入。一个商品实时推荐系统,行为流和静态数据是两条路。

数据类型推荐系统里的用途推荐接入方式
用户行为日志计算点击率、实时偏好Kafka JSON 流
商品属性与价格召回阶段的品类匹配MySQL CDC 或定时同步
用户画像特征拼接Redis 读取

行为日志走 Kafka 是业界主流,压缩包里的自定义 DataSource 本质上是模拟 Kafka 消息。如果项目里出现 MySQL CDC,通常是为了同步商品表,而不是同步行为日志。两者职责不同,别混在一起。

Flink CDC 安装部署和 pipeline 部署最近讨论多,但放到推荐项目里,我建议控制使用范围。商品表的变更频率低、数据量小,用 Flink CDC 的initial模式,一次性快照加后续 binlog 增量,性价比最高。有些新手把整张订单表都接入 CDC 再过滤行为,这会拖慢作业,还会让 MySQL 压力陡增。

CREATE TABLE products ( item_id INT PRIMARY KEY, category_id INT, price DECIMAL(10, 2), tags STRING ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '127.0.0.1', 'port' = '3306', 'username' = 'flinkuser', 'password' = 'flinkpw', 'database-name' = 'shop', 'table-name' = 'products', 'scan.startup.mode' = 'initial' );

这个 SQL 里的scan.startup.mode建议别默认。initial模式适合需要完整快照的场景,任务重启时会重新做一次快照,表大时相当耗时。如果商品的当前状态就能满足推荐需要,用latest-offset会快很多。JDBC 连接器在这里也被经常提,它和 CDC 连接器是两回事,JDBC 用于读取或写入维表,CDC 用于监听变更日志,第 4.1 节会单独讲 JDBC 的坑。

3. 实时推荐的三个核心算子:窗口聚合、状态去重与 Redis 特征回写

链路通了,接着是把推荐特征计算出来。推荐系统里常见的做法是分三层:行为计数、去重、特征存储。这一章的三个小节分别对应这三个环节,代码可以直接抄,但参数要按你的业务调。

3.1 滑动窗口做行为计数:从词频统计到曝光/点击特征

Flink 入门都会做词频统计,很多人觉得那只是个 warm-up,其实推荐系统的行为计数和它一模一样。你现在统计的是用户过去 10 分钟内点击了某个商品多少次,就是 WordCount 换了一层皮。

DataStream<UserBehaviorCount> countStream = stream .filter(e -> "click".equals(e.getBehavior())) .keyBy(e -> e.getUserId() + ":" + e.getItemId()) .window(SlidingProcessingTimeWindows.of(Time.minutes(10), Time.minutes(5))) .aggregate(new AggregateFunction<BehaviorEvent, long[], UserBehaviorCount>() { @Override public long[] createAccumulator() { return new long[1]; } @Override public long[] add(BehaviorEvent event, long[] acc) { acc[0] += 1L; return acc; } @Override public UserBehaviorCount getResult(long[] acc) { return new UserBehaviorCount(acc[0]); } @Override public long[] merge(long[] a, long[] b) { a[0] += b[0]; return a; } });

滑动窗口的步长决定特征更新的频率。10 分钟窗口配 5 分钟步长,意味着每个用户商品对每 5 分钟输出一版新计数。步长越短,特征越实时,但窗口计算和下游存储的压力也越大,这个读写比要心里有数。

这里我特意用AggregateFunction而不是ProcessWindowFunction。原因是AggregateFunction是增量聚合,窗口数据只保留一个 accumulator,内存占用恒定;ProcessWindowFunction要缓存整个窗口的全部数据,商品量一大,JVM 堆直接告急。窗口聚合时,createAccumulator和getResult的触发机制可以这样理解:每条数据来,add更新累加器;窗口触发,getResult取结果。它的性能比全量窗口高一个数量级。

生产环境别用SlidingProcessingTimeWindows,用事件时间和水印来控制数据迟到的表现。设置水印的方式是env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)在旧版本中常见,Flink 1.13 之后推荐直接用WatermarkStrategy,否则乱序数据会直接造成推荐特征偏差。

3.2 用 Keyed State 做去重:TTL 参数这样设

推荐场景里有个常见需求:同一用户一天内对同一商品的曝光只算一次,或者点击去重。实时流里最简单的去重是拿 HashSet 放在内存里,但分布式环境下这是黑匣子,TaskManager 一重启就全丢。正确做法是用 Keyed State。

public class DeduplicateFunction extends KeyedProcessFunction<String, BehaviorEvent, BehaviorEvent> { private ValueState<Boolean> seenState; @Override public void open(Configuration parameters) { ValueStateDescriptor<Boolean> descriptor = new ValueStateDescriptor<>("seen", Types.BOOLEAN); StateTtlConfig ttl = StateTtlConfig.newBuilder(Time.hours(24)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); descriptor.enableTimeToLive(ttl); seenState = getRuntimeContext().getState(descriptor); } @Override public void processElement(BehaviorEvent value, Context ctx, Collector<BehaviorEvent> out) throws Exception { String dedupKey = value.getUserId() + ":" + value.getItemId(); Boolean seen = seenState.value(); if (seen == null) { seenState.update(true); out.collect(value); } } }

这里的StateTtlConfig是去重的生命周期闸门。24 小时的 TTL 意味着用户对同一个商品的去重只持续一天,第二天状态自动失效,新的行为重新计数。OnCreateAndWrite表示每次对状态做写操作时都会刷新过期时间,适用于连续活跃的用户;如果你希望过期时间严格从创建时间算,哪怕用户一直在产生行为也要过过期时间,那就用OnCreateAndWrite之外的选项,具体行为要以 Flink 版本自带的枚举为准。

NeverReturnExpired值得展开说。它保证即使状态还没被后台清理,读取时也绝不返回已过期的数据。反过来ReturnExpiredIfNotCleanedUp可能让你读到一秒前就已过期的数据,造成推荐重复。去重逻辑选前者,没得商量。

状态 TTL 还有个隐蔽问题。StateTtlConfig过期数据要等后台清理线程触发才会真正从存储中删掉,如果状态本身不常被访问,它就一直是僵尸数据,RocksDB 体积只增不减。此类问题通常可以通过周期性访问或者设置过期状态的清理策略来控制,具体设置项不同版本差异较大,我一般是先让作业跑一两天,观察 checkpoint 大小的增长曲线再决定要不要调整。

3.3 特征回写 Redis:自定义 DataSink 的参数与序列化

特征算完不能放在 Flink 状态里裸着,推荐服务要主动读。业界最常见的是把特征写回 Redis,商品 ID 做 key,多个特征做 hash field。下面的 Sink 是 RichSinkFunction 的典型实现。

public class FeatureRedisSink extends RichSinkFunction<FeatureRow> { private transient JedisPool pool; private String redisHost; private int redisPort; private int maxTotal = 16; public FeatureRedisSink(String redisHost, int redisPort) { this.redisHost = redisHost; this.redisPort = redisPort; } @Override public void open(Configuration parameters) { JedisPoolConfig config = new JedisPoolConfig(); config.setMaxTotal(maxTotal); config.setMaxWaitMillis(3000); config.setTestOnBorrow(true); pool = new JedisPool(config, redisHost, redisPort, 3000); } @Override public void invoke(FeatureRow row, Context context) throws Exception { try (Jedis jedis = pool.getResource()) { String key = "rec:item:" + row.getItemId(); jedis.hset(key, row.getFeatureName(), row.getFeatureValue()); jedis.expire(key, 24 * 3600); } } @Override public void close() throws Exception { if (pool != null) { pool.close(); } } }

open里创建连接池,invoke里每次获取连接并写数据。这里有个玄学或者说血泪经验:连接池的maxTotal不要设太大,16 到 32 足够。Flink 的 Sink 算子并行度若开到 16,每个子任务都有自己的连接池,总连接数就是 16 乘以 16,Redis 默认 maxclients 会被打满,现象就是下游服务间歇性超时。

为什么用hset而不是set?因为一个商品的推荐特征往往有多个维度,点击量、加购量、类别标签,放同一个 key 的不同 field 里,推荐服务一次 HGETALL 就能拿到全量特征,避免多次网络往返。那个expire(key, 24 * 3600)是保险丝,防止有些商品长期没有新行为,Redis 里残留旧特征误导推荐排序。

生产环境我会在这里做序列化改造。上面的代码直接把特征值写成字符串,特征是数字没问题,如果特征是 flag 或 embedding,就要考虑压缩或二进制格式。你可以在invoke里加个简单的判断,比如长字符串字段走 Protobuf 或 Kryo,数值字段继续走字符串。这是我踩完坑之后的习惯,别把整行数据塞进 Redis 当 RedisJSON 用,查询和内存代价都不低。

4. 实时推荐项目避坑指南:JDBC 连接器异常、Sink Hive 不入表与背压排查

这一章是硬菜。实时推荐系统翻车,大概率不是算法不对,而是连接器、存储和资源这三类问题。我把高频坑按现象、原因、解决三段式写清楚,你可以直接拿来当排查手册。

4.1 JDBC 连接器异常:先查驱动、再查时区、最后查连接复用

先说最常见的现象:任务提交到集群时直接挂掉,报错信息里有ClassNotFoundException: com.mysql.cj.jdbc.Driver,或者Communications link failure。

驱动找不到的原因绝大多数是 Jar 包没带驱动。Flink 官方 lib 目录里不会有 MySQL 驱动,它只带基础的连接器和核心依赖。解决方式是把驱动打进作业 Jar,并检查最终产物。

<dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>8.0.33</version> <scope>compile</scope> </dependency>
unzip -l target/recommend-1.0.jar | grep -i mysql

第一条是 maven 依赖,第二条命令检查最终 Jar 里有没有 driver。如果你用 shade 插件打包时配置了filters排除掉 META-INF 下文件,可能导致驱动服务的 SPI 加载失效,最直接的现象是驱动类找到了但初始化失败。这个坑我踩过一次,最后把ServicesResourceTransformer加进 shade 插件,问题才消失。

再说Communications link failure。这个报错常见于作业跑了几小时后出现,原因和 Flink 自身的 JDBC 连接器设计有关:连接器会维护一个连接池,空闲时间超过了 MySQL 的wait_timeout,MySQL 那边把连接断了,而 Flink 这边还在复用这条死连接。解决方式我推荐两个方向:

参数作用建议值
sink.buffer-flush.max-rows攒够多少条开始批量写100 到 1000
sink.buffer-flush.interval多久强制刷一次缓冲1 到 3 秒

这两个参数是 JDBC sink 的核心。sink.buffer-flush.interval设得太长,推荐结果延迟高;太短,连接频繁建立反而触发 MySQL 端限流。还有一种邪道做法,定期重建连接池的定时任务,虽然能绕开wait_timeout,但引入的复杂度更高,不建议在入门项目里用。

4.2 Sink Hive 表数据不入表:先查缓冲与分区提交

另一个高频事故:作业显示正常结束、记录数也对,但去 Hive 查表没有新数据。这个问题在流式写 Hive 时特别典型。

现象是双重的:一是 Hive 表在查询时看不到数据,二是 Flink 任务端的指标显示已经写入了 1 万多条记录。原因往往是写入数据先被缓冲在 Orc/Parquet writer 的内存里,没达到触发文件切分的条件,文件一直没生成,Hive 表自然查不到。

需要检查的第一类是文件滚动参数:

参数作用推荐值
sink.rolling-policy.file-size文件大小触发滚动128MB
sink.rolling-policy.rollover-interval时间间隔触发滚动30 分钟
sink.partition-commit.trigger分区提交时机partition-time
sink.partition-commit.policy.kind提交动作metastore, success-file

当你发现 Hive 表里数据迟迟不出现时,用这两个思路排查:先看sink.rolling-policy.rollover-interval有没有设置,把时间改为 5 分钟,能快速看到新文件生成;然后看sink.partition-commit.trigger是不是 partition-time,如果是,需要检查分区提取的时间模式对不对,常见问题是你写的时间戳是yyyy-MM-dd HH:mm:ss,但表分区是dt格式,提取器拿不到正确的时间,分区永远不提交。

还有一个和这两个无关的场景:如果你写的是非分区表,Hive 表能看到数据但落在临时目录里,那是 Flink 的 filesystem sink 把文件先写到了 staging 目录,需要 job 关闭 commit 策略确认后才会 rename 到正式目录,这种情况不是相关性的 bug,是提交策略还没配置好的表现。

4.3 背压升高与火焰图定位:排查顺序别搞反

任务越跑越慢,Kafka 消费 lag 持续增长,这就是背压的典型信号。推荐系统对延迟敏感,背压一旦出现,用户点击后特征更新滞后,推荐结果完全失去实时意义。

先打开 Web UI 的 Backpressure 页。你会看到每个算子的背压状态,分 OK、LOW、HIGH 三档。排查时别一上来就抓线程栈,先看背压状态在哪个算子处发生。如果某个 keyBy 下游所有子任务都是 HIGH,那是典型的数据倾斜;如果所有子任务都 HIGH,且只出现在 Sink 算子,优先怀疑下游存储瓶颈。

定位到具体算子后,再上火焰图。Flink 的火焰图一般有两种获取方式:一是 Web UI 里直接触发栈采样,另外就是自己在 TaskManager 进程上用命令抓。

jps -l | grep TaskManagerRunner jstack <pid> > jstack_$(date +%s).txt

等 30 秒再执行一次,两次栈采样对比,你能看到某个方法反复出现在栈顶。比如RedisSink.invoke占了大半,说明写 Redis 慢,回到 3.3 节检查序列化和连接池配置;如果是KeyedProcessFunction的定时器处理占大头,说明注册的定时器太多,去重或窗口过期的状态清理机制有隐患。

特别提醒,别把背压和延迟画等号。Kafka lag 升高也可能是 source 消费能力不足,但背压显示 OK,这里要查分区数是 Kafka 侧的问题,不是 Flink 的问题。排查顺序固定为:先确认背压出现在哪个算子,再上采样定位具体瓶颈,最后才调并行度和参数。顺序反了会乱改一通,越改越慢。

4.4 Checkpoint 恢复的坑:状态丢、UID 变、路径漂

实时推荐的恢复能力靠 Checkpoint 兜底。我见过最伤的一次事件:凌晨集群扩容,作业重启,结果用户近一天的行为状态全丢,推荐特征重新积累,花了 6 个小时才恢复正常。原因是当时 checkpoints 目录配的是本地路径,TaskManager 重启后文件直接被清掉。

flink run -s hdfs://nameservice/flink/checkpoints/recommend-job/_metadata \ -c com.recommend.job.RecommendationJob \ -p 4 \ target/recommend-1.0.jar

-s参数指定从某个 checkpoint 元数据恢复。生产环境 checkpoints.dir 一定要放在共享存储,本地文件系统只能用在单机演示。另一个恢复时的高发问题,是改代码后状态对不上。

状态的 Key 是seen、窗口聚合算子的状态 ID 也是自动生成的。只要你在代码里加了一行算子,flink 重算生成的 uid 就变了,恢复时找不到旧状态,表现就是任务起来了但状态为空。解决方式是给关键算子显式设置 uid:

stream .keyBy(e -> e.getUserId()) .window(...) .uid("user-click-window") .name("user-click-window");

uid("user-click-window")固定了算子身份,即使代码顺序变化,只要 uid 不变,状态依然匹配。这个习惯要早点养成,等状态已经恢复不上再改 uid,那才叫真正的翻车。还有一个常被人忽略的点:RocksDB 状态恢复时要反序列化旧数据,如果状态里存的对象的 Serializable 版本变了,会直接报反序列化异常。因此,状态里存的对象尽量保持简单且不轻易删字段,宁可新增字段,也别改已有字段的类型。

5. 进阶验证:离线回放、影子流量与血缘管理收口

项目跑通只是起点,怎么证明推荐结果是好的,这题很多人没想清楚。我先给一个低成本验证方案:把真实日志按时间回放一遍。做法是把 Kafka 里的历史消息保留两天,用一个独立 group 让 Flink 作业从最早 offset 开始重新消费,不回头的所有效果都用同一套参数跑一遍。你会在几分钟内看到之前真实用户的行为如何映射为推荐的候选集,比拿着代码空想实在得多。

回放之前,我建议先在代码里加一个流量切分逻辑,用影子流量方式对比新旧两套特征效果。比如按用户 ID 哈希取模,10% 的用户走新版本特征:

# 这是一段贴近线上业务逻辑的简化示意,具体用 Java 还是引擎表达式按团队来 def should_eval_new_version(user_id: str) -> bool: return hash(user_id) % 100 < 10

影子流量的意义在于,不让新算法影响全部真实用户,只让新版本在副路径上跑,输出写到影子 Redis 前缀,业务方仍然读主版本的 Redis。对比时统计两边的点击纵坐标升幅,才能证明新特征的收益。我做这类验证时,最常看反直觉的结论:新特征在离线回测里表现好,线上却可能完全死掉,原因往往不是特征本身,而是延迟稳定性。

最后收口做血缘管理。项目一多,别人拿着一份特征数据问你“这是从哪张表来的、哪天开始生效”,如果你答不上来,后续排查和生产事故定位就会很被动。生产里我见过用 OpenMetadata 抓取 Flink 血缘关系的团队,也有用数据字典硬记的做法。我的实际建议是:不管上不上元数据平台,至少给每个 Flink 作业起名时把主题、表、版本都带上,比如feature_user_ctr_v3,别让作业名变成默认的 jar 名。

这是我吃了一次亏换来的教训:上线前的回放、切换时的影子流量、上线后的血缘记录,这三件事一步都不能省。希望帮到你。

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

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

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

立即咨询