
简介基于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/checkpointstaskmanager.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 RichSourceFunctionBehaviorEvent { 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(SourceContextBehaviorEvent 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)控制生成速率这一行非常关键它决定了你的窗口计算能不能观察到滑动效果。用的时候这样加进作业DataStreamBehaviorEvent 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 换了一层皮。DataStreamUserBehaviorCount 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 AggregateFunctionBehaviorEvent, 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 KeyedProcessFunctionString, BehaviorEvent, BehaviorEvent { private ValueStateBoolean seenState; Override public void open(Configuration parameters) { ValueStateDescriptorBoolean 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, CollectorBehaviorEvent 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 RichSinkFunctionFeatureRow { 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 乘以 16Redis 默认 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 groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version8.0.33/version scopecompile/scope /dependencyunzip -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_timeoutMySQL 那边把连接断了而 Flink 这边还在复用这条死连接。解决方式我推荐两个方向参数作用建议值sink.buffer-flush.max-rows攒够多少条开始批量写100 到 1000sink.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文件大小触发滚动128MBsink.rolling-policy.rollover-interval时间间隔触发滚动30 分钟sink.partition-commit.trigger分区提交时机partition-timesink.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 就变了恢复时找不到旧状态表现就是任务起来了但状态为空。解决方式是给关键算子显式设置 uidstream .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 名。这是我吃了一次亏换来的教训上线前的回放、切换时的影子流量、上线后的血缘记录这三件事一步都不能省。希望帮到你。本文还有配套的精品资源点击获取