
简介一份围绕Flink实时计算与推荐算法构建的商品实时推荐系统项目适合大数据方向学生、初级工程师以及希望上手实时推荐场景的开发者参考。项目主体为34个Scala源码文件覆盖数据接入、窗口统计、特征处理与推荐结果输出等环节并包含2个SQL脚本、2个HBase建表语句、2个properties配置及Kafka模拟数据脚本便于从环境搭建到数据流调试完整跑通。压缩包共44个文件整体仅245KB轻量但结构清晰适合快速阅读与二次开发。目前已有275人学习下载。通过阅读工程代码可掌握Flink DataStream API在用户行为实时分析中的应用思路理解协同过滤等推荐逻辑的落地方式同时参考其状态管理、Kafka与HBase集成等细节为构建个性化推荐服务打下基础。1. 基于 Flink 的商品实时推荐这份资源到底能帮你落地什么做大数据实时计算的人绕不开 Flink做电商推荐系统的人绕不开实时用户行为。这份《基于 Flink 商品实时推荐系统.zip》把两者合在了一起——里面是一套完整的 Flink 推荐系统工程源码主目录是flink-recommend-system-main能看到flink-2-hbase、pom.xml、src、data、sql等标准工程结构。它不是某个课程里摘出来的零散片段而是从数据接入、预处理、特征工程、推荐算法到结果落库的一个可运行闭环。对正在学 Flink 或者准备做毕设、做面试项目的人来说这份资源直接给你一个能拆、能跑、能改的参照系。它能解决的核心问题很具体用户点击、浏览、购买这些行为数据来了之后Flink 怎么在秒级窗口内算出「这个人接下来最可能买什么」。适合三类人刚学完 Flink 基础想搞懂流处理怎么用在真实场景的正在设计推荐系统架构需要参考代码的以及想把数据从 Kafka 或日志文件实时写到 HBase 再做查询的。2. 架构选型与工程结构先看清 Flink 在推荐链路里站在哪个位置2.1 实时推荐系统的整体数据流从行为日志到推荐结果的四段式链路一份合格的 Flink 工程拿到手第一件事不是看代码而是看数据是怎么流起来的。这套系统的数据流可以拆成四个阶段数据源接入、预处理与宽表构建、特征统计与算法计算、结果输出。数据源阶段用户的点击、曝光、加购、下单行为来自前端埋点日志或消息队列在这套工程里data目录下放了模拟行为数据src里的 Source 实现会按固定速率读取这些数据模拟真实日志流。预处理阶段用 DataStream API 做过滤、去重、字段补齐把乱七杂八的原始日志转成统一的用户行为事件对象。特征工程阶段是这套系统的重头戏通过窗口函数计算用户最近 N 分钟的浏览频次、商品类目偏好等特征。算法阶段则承担了协同过滤或基于规则的实时召回。最后结果输出到 HBase也就是flink-2-hbase这个模块干的事下游的推荐接口或者 Web 端直接查 HBase 就能拿到推荐列表。从工程角度看这个架构最值得学习的地方在于它把「实时计算」和「在线存储」解耦了。Flink 不直接面向用户提供服务而是把计算好的推荐结果写入 HBase由查询端按需读取。这样做的好处是推荐计算和用户请求不互相阻塞——就算 Flink 任务重启或者做检查点恢复用户端依然能读到上一次写入的推荐结果只是延迟更新而已。这种设计在实际生产里非常常见也是面试官愿意听到的答案实时链路要保证吞吐在线链路要保证低延迟两者之间用一个支持随机读的 KV 存储来衔接。2.2 工程目录与模块职责pom.xml、src、data、sql 分别管什么打开压缩包你会看到几个鲜明的顶层节点。pom.xml是 Maven 多模块工程的父 POM里面定义了 Flink 版本、依赖管理、插件配置。flink-2-hbase是一个独立 Maven 模块从名字看就是 Flink 计算后写入 HBase 的落地模块里面包含了 HBase 的 Sink 实现和表结构映射。src是主代码目录按 Flink 标准结构分为 main 和 testmain 下通常有java和resourcesJava 代码里包含了 Source、ProcessFunction、窗口计算、UDF、Sink 等核心逻辑。data目录放着模拟的用户行为数据文件一般是 CSV 或者 JSON 格式字段包含用户 ID、商品 ID、行为类型、时间戳等。sql目录则存放了 HBase 建表语句或者 Flink SQL 的初始化脚本比如你要先建好 HBase 的命名空间和表再跑 Flink 作业。我建议拿到资源后按这样的顺序去读代码先看pom.xml里 Flink 的版本和依赖范围再看sql目录确认外部存储的表结构接着读src/main/java下的主入口类顺着 DataStream 的处理链往下走。不要一上来就钻进具体的函数实现先画清楚输入和输出的接口再往里面填细节。这份工程里有一个典型的模式——主类里先env.addSource()得到 DataStream然后接.filter().map().keyBy().window().process()最后.addSink()落库你只要找那串链式调用整个系统的主干就出来了。2.3 为什么选择 HBase 作为推荐结果存储随机读与批量写入的平衡项目把结果写到 HBase 而不是 MySQL这个选择是有讲究的。推荐结果的查询模式是「给定用户 ID返回商品列表」这是典型的点查场景HBase 的 RowKey 设计正好能支持。如果 RowKey 是userId的哈希前缀那一次查询就能精确命中一个 Region毫秒级返回。写方面Flink 的实时计算结果是一个持续不断的流每条数据都需要写入HBase 的批量写入接口BufferedMutator能攒一批再发避免了逐条插入的网络开销。相比之下MySQL 在持续写入和随机读并发同时走高时容易成为瓶颈而且 schema 固定很难灵活存储每个用户不同长度的推荐列表。另一个原因是列族的灵活性。HBase 可以在一个 RowKey 下存储多个列比如recs:item1、recs:item2也可以存成一个 JSON 字符串放在一列里。这套系统的实现里我猜大概率是把推荐的商品 ID 列表序列化后存进一列这样读写最简单。如果你要改造自己的项目我建议保留这个方案——不要为推荐结果设计太复杂的表结构推荐列表本质上是「一个用户对应一串商品」用一个列存 JSON 完全够用。3. 数据接入与预处理把原始行为日志变成干净的事件流3.1 Source 实现与模拟数据读取为什么用 DataStream 而不是 Flink SQL在 Flink 1.12 之后Flink SQL 开始被大量使用但这份工程仍然走的是 DataStream API。原因是推荐系统里有很多面向单个事件的精确操作——比如判断行为类型、修正时间戳、拼接上下文信息用 SQL 的声明式语法表达反而绕。DataStream 的算子链更直白每个 map、flatMap、process 里写的就是一段 Java 逻辑调试时打日志方便断点也好下。对学习而言DataStream 能让你把流处理的每个环节看得明明白白不会被 SQL 的优化器遮住细节。项目里的数据源很可能是从data目录读取文件的常见做法是用自定义 SourceFunction 循环读取文件内容或者用readTextFile加env.setParallelism(1)控制读取。如果是真实场景这里应该替换成 Kafka Source比如DataStreamString kafkaStream env.addSource( new FlinkKafkaConsumer010(user-behavior, new SimpleStringSchema(), kafkaProps) );逻辑说明addSource方法接收一个 SourceFunction这里用的是 FlinkKafkaConsumer010它会把 Kafka 里user-behavior主题的每条消息以字符串形式传给下游。参数方面new SimpleStringSchema()指定反序列化方式kafkaProps里要配bootstrap.servers、group.id、auto.offset.reset等。如果你用的是新版本 Flink请改用KafkaSource.builder()那一套 API老接口虽然还能跑但官方已经标注过期。3.2 清洗与转换去重、异常值过滤、时间戳标准化原始日志永远是脏的。常见的垃圾数据有重复上报的点击事件、时间戳乱序、用户 ID 为空、商品 ID 不在商品表里。这一阶段用 DataStream API 逐个算子去处理。先做过滤把明显无效的事件扔掉DataStreamUserBehavior cleaned rawStream .map(new MapFunctionString, UserBehavior() { Override public UserBehavior map(String line) throws Exception { String[] fields line.split(,); return new UserBehavior( Long.parseLong(fields[0]), Long.parseLong(fields[1]), fields[2], Long.parseLong(fields[3]) * 1000L // 秒转毫秒 ); } }) .filter(behavior - behavior.userId 0 behavior.itemId 0) .filter(behavior - behavior.timestamp 0);逻辑说明第一个 map 把一行 CSV 文本切分并封装成 UserBehavior POJO第三个字段是行为类型比如 view、cart、buy第四字段时间戳在数据源里可能是秒这里乘 1000 转成毫秒方便和 Flink 的 EventTime 对齐。后面两个 filter 分别过滤掉非法用户和非法商品以及时间戳为负的异常数据。参数注意时间戳单位不统一是新手最容易翻车的地方如果数据源里已经是毫秒你还乘 1000时间就膨胀了 1000 倍窗口全乱。3.3 Watermark 与乱序处理实时推荐里被忽略的根基实时推荐最怕的不是慢是乱序。用户的行为事件从不同设备、不同网络环境上报到服务器到达 Kafka 或 Flink 的时间不一定和事件发生时间一致。如果你用 ProcessingTime处理时间来开窗口那么窗口边界就是「数据到达 Flink 的时刻」这会导致一个问题一个用户先浏览了 A 商品网络延迟导致这条事件后到系统会把它算进下一个窗口特征就错了。正确做法是用 EventTime配合 Watermark 来处理乱序。DataStreamUserBehavior withWatermark cleaned .assignTimestampsAndWatermarks( WatermarkStrategy.UserBehaviorforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, timestamp) - event.timestamp) );逻辑说明forBoundedOutOfOrderness(Duration.ofSeconds(10))表示允许数据最多晚到 10 秒晚于这个界限的数据会被丢弃或进入侧输出流。withTimestampAssigner告诉 Flink 从哪个字段提取事件时间。参数调整建议这个 10 秒不是拍脑袋定的要看上游数据的延迟分布电商日志 95% 的事件在 5 秒内到达那就设 6 到 8 秒如果你发现某段时间推荐明显不准先看监控里 late data 的比例如果很高说明 Watermark 设得太大窗口太晚触发推荐结果时效性差反之如果 late data 接近 0说明你牺牲了太多延迟来等数据。实时推荐系统里等待和准确永远是矛盾的我一般先设一个偏小的值跑几天看延迟报告的曲线再微调。4. 特征工程与窗口计算从行为流中提取用户偏好4.1 基于时间窗口的统计特征最近 N 分钟浏览频次与商品类目偏好推荐系统里最基础也最有效的特征之一是用户最近一段时间内的行为频次。比如「最近 5 分钟看了哪些商品、看了多少次」「最近 30 分钟加购了几个商品」。这些特征能迅速反映用户的即时兴趣——一个人刚搜了「机械键盘」系统就应该在这个时间窗口内给他推键盘配件。Flink 的窗口计算是干这个的天然工具。滑动窗口能保证每隔几秒更新一次推荐结果让我们拿SlidingEventTimeWindow来举例DataStreamUserBehavior windowed withWatermark .keyBy(behavior - behavior.userId) .window(SlidingEventTimeWindows.of( Time.minutes(10), // 窗口大小 Time.minutes(1))) // 滑动间隔 .process(new CountAggregateFunction());逻辑说明keyBy按用户 ID 分组每个用户的独立事件流各自开窗。of的第一个参数是窗口长度 10 分钟第二个参数是滑动步长 1 分钟也就是说每隔 1 分钟输出一次最近 10 分钟的统计值。process接收一个自定义 ProcessWindowFunction在里面可以拿到窗口内所有元素做计数、去重、累加之类的操作。参数设计中窗口越大统计越平滑但对瞬时兴趣反应越慢滑动越频繁则计算开销越大。对于电商实时推荐我习惯把窗口控制在 5 到 15 分钟之间滑动间隔不超过窗口长度的四分之一这样既能捕捉短期兴趣又不会让状态无限膨胀。4.2 用 ProcessFunction 实现用户行为序列把离散事件拼成连续轨迹窗口统计能给数值特征但用户的行为顺序同样重要。比如「先点了 A 商品又点了 B 商品最后加购了 C」——这个序列里隐含了强关联信息能为协同过滤提供高质量的输入。Flink 没有内置的序列拼接算子但 ProcessFunction 可以把每个用户窗口内的行为按时间排序后拼成一个字符串或列表public class BehaviorSequenceFunction extends ProcessFunctionUserBehavior, UserSequence { private ValueStateListUserBehavior sequenceState; Override public void open(Configuration parameters) { ListStateDescriptorUserBehavior descriptor new ListStateDescriptor(behavior-seq, UserBehavior.class); sequenceState getRuntimeContext().getListState(descriptor); } Override public void processElement(UserBehavior value, Context ctx, CollectorUserSequence out) throws Exception { long currentWatermark ctx.timerService().currentWatermark(); // 把当前事件加入列表状态 sequenceState.add(value); // 注册一个 30 秒后的定时器统一输出 ctx.timerService().registerEventTimeTimer(currentWatermark 30000); } Override public void onTimer(long timestamp, OnTimerContext ctx, CollectorUserSequence out) throws Exception { IterableUserBehavior behaviors sequenceState.get(); ListUserBehavior sorted new ArrayList(); for (UserBehavior b : behaviors) { sorted.add(b); } sorted.sort(Comparator.comparingLong(b - b.timestamp)); // 拼成 itemId1:view,itemId2:view,itemId3:cart 这样的字符串 StringBuilder sb new StringBuilder(); for (UserBehavior b : sorted) { sb.append(b.itemId).append(:).append(b.action).append(,); } if (sb.length() 0) { out.collect(new UserSequence(sorted.get(0).userId, sb.toString())); } sequenceState.clear(); } }逻辑说明ListState用于保存该 key 下累积的行为列表这里 keyBy 的用户 ID。每次来一条数据先存进状态然后注册一个基于事件时间的定时器到了触发时刻就把所有行为取出来按时间排序拼接成序列字符串输出。参数说明定时器延迟 30 秒是为了等一批行为积攒后再输出减少下游处理频率如果你对实时性要求高可以改成 10 秒但要注意定时器过多会加重 Flink 的定时器状态负担。这里的核心思想是「用状态攒批用定时器控制发射节奏」比每个事件都直接下发更高效。4.3 状态管理为什么统计频次要用 Flink State 而不是外部 Redis很多初学者会想统计用户频次直接每次去 Redis 里 INCR 不行吗在演示项目里行在真实生产里不行。原因有两个一是网络开销每来一条事件就访问一次 Redis吞吐上去了 Redis 压力极大而且多了一次 RTT二是无法和 Flink 的容错机制协同——如果 Flink 任务崩溃Redis 里的计数和 Flink 的状态可能不一致恢复后你无法判断哪些数据已经算过。Flink 的 Keyed State 由 RocksDB 或内存管理随检查点持久化。用上面代码里的ValueState、ListState或更简单的ReducingState来做累加任务失败后检查点会自动恢复现场。代价是你得把「调外部存储记账」的思维转成「flink 是单机状态的延续」——状态天然和 key 绑定绝不混用。若你真的想把计数结果暴露给外部查询正确做法是等窗口计算完成后再写入 HBase 或 Redis而不是在计算过程中实时读写外部系统。5. 推荐算法与结果输出把特征变成最终的商品列表5.1 实时推荐算法的选择为什么拿 Item-Based 协同过滤做主力实时的场景决定了算法不能太复杂。矩阵分解、深度学习模型在离线批量推荐里效果好但它们需要训练过程模型更新周期长很难在秒级内响应一个新变化。Item-Based 协同过滤物品协同过滤是实时推荐里最务实的选项算法预计算好商品之间的相似度矩阵这份相似度可以离线算好存在 HBase 或本地文件里在线阶段只需要获取用户最近交互过的商品去查这些商品的相似商品按相似度加权排序就得到推荐列表。这个过程的计算量只和用户最近交互的商品数有关和全量商品数无关非常适合在 Flink 的流处理里实时完成。这种做法有两个好处一是不依赖用户历史长序列用户刚点了一个新商品马上能推荐它最相似的商品冷启动能力强二是相似度矩阵可以定时更新比如每天凌晨用离线作业重算一次白天 FLink 只做查表和排序压力小。如果你要在这份工程基础上升级可以先用离线统计的方式把「商品共现矩阵」算出来再在 Flink 里读取替换掉原来可能内置的简单规则逻辑。5.2 用 HBase 存储相似度矩阵与推荐结果的读写实现既然 Final 结果要写 HBase那就用 HBase 做相似度矩阵的载体RowKey 是商品 ID列族sim:itemId2存相似度值。Flink 侧用异步 I/O 查 HBase避免同步查阻塞流计算。不过这份工程里更简单的做法是直接把离线算好的相似度表加载为 Flink 的广播状态BroadcastState每个算子都持有一份只读副本本地查表完全无网络开销。广播状态适合数据量不大几万商品的场景商品过百万就该换 HBase 了。输出推荐结果时通常用 Tuple 或 POJO 封装用户 ID 和商品列表写入 HBase。RowKey 设计成userId的哈希加时间戳前缀避免热点列存储推荐商品 JSON 串如下// 推荐结果写入 HBase public class HBaseSink extends RichSinkFunctionRecommendResult { private Connection connection; private BufferedMutator mutator; Override public void open(Configuration parameters) throws Exception { ConnectionFactory factory new ConnectionFactory(); connection factory.createConnection(simpleConfig); BufferedMutatorParams params new BufferedMutatorParams(TableName.valueOf(rec:result)) .writeBufferSize(8 * 1024 * 1024); // 8MB 缓冲 mutator connection.getBufferedMutator(params); } Override public void invoke(RecommendResult value, Context context) throws Exception { String rowKey value.userId _ (System.currentTimeMillis() / 60000); Put put new Put(Bytes.toBytes(rowKey)); put.addColumn(Bytes.toBytes(rec), Bytes.toBytes(items), Bytes.toBytes(value.itemListJson)); mutator.mutate(put); } Override public void close() throws Exception { mutator.flush(); mutator.close(); connection.close(); } }逻辑说明open里创建 HBase 连接和 BufferedMutator8MB 缓冲让 Flink 攒够一批再写入避免每条结果都做一次 RPC。invoke里每条结果构造一个 PutRowKey 是用户 ID 加分钟级时间戳列存 JSON。close里 flush 清空剩余缓冲。参数注意写缓冲不能设太大否则容器内存会涨也不能太小否则退化成逐条写。推荐 4 到 16MB要是你发现写入耗时波动大看下区域服务器 GC很可能是缓冲配置和 Region 分裂策略不匹配。5.3 用 Flink CDC 实时同步 MySQL 到 ClickHouse一份工程外的扩展这套工程本身不包含 ClickHouse但在实际部署里「使用 Flink 实现 MySQL 同步到 ClickHouse」是一个高频搭配——推荐结果从 HBase 读出后运营分析团队可能希望把点击日志同步进 ClickHouse 做离线多维分析。常见做法是用 Flink CDCChange Data Capture监听 MySQL 的 binlog把变更实时写入 Kafka再从 Kafka 用 Flink 消费写入 ClickHouse-- Flink SQL 里创建 CDC 源表和 ClickHouse 目标表 CREATE TABLE mysql_orders ( id INT, user_id BIGINT, item_id BIGINT, action STRING, -- 必须声明主键否则 Flink CDC 无法识别 binlog 主键 PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username cdc_user, password ******, database-name shop, table-name orders ); CREATE TABLE clickhouse_orders ( id INT, user_id BIGINT, item_id BIGINT, action STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector clickhouse, url clickhouse://localhost:8123, table-name orders ); INSERT INTO clickhouse_orders SELECT * FROM mysql_orders;逻辑说明mysql-cdc连接器自动捕获 binlog 的插入、更新、删除事件并在 Flink 内转换为对应的 changelog 数据流clickhouse连接器接收 changelog 并同步到 ClickHouse 的 MergeTree 表。参数注意MySQL 必须开启 binlog_formatROW连接用户需要 REPLICATION SLAVE 权限ClickHouse 表必须有主键对应 CDC 的 upsert 语义否则重复数据会叠加而不是更新。这一步做到推荐系统外面能让你的整体大数据链路完整度提高不少——流上算了什么、实时同步了什么都能对得上。6. 避坑与排查Flink 实时推荐常见的六个实战雷区6.1 现象Flink 作业启动后一直报「flink 的 jdbc 连接器异常」原因新版 Flink 把 JDBC 连接器从默认依赖改成可选你需要显式引入flink-connector-jdbc而且版本要与 Flink 主版本匹配。另一个原因是 HBase 连接器同样如此——flink-2-hbase模块如果没用对 HBase 客户端版本会报类冲突。解决在pom.xml里显式声明连接器依赖并剔除自带的老版本客户端。比如 Flink 1.14 配 HBase 2.x通常用org.apache.flink:flink-connector-hbase-2.2。升级版本时注意 Flink 与连接器的编译兼容矩阵不是越新越好而是要和主版本一致。6.2 现象窗口聚合结果迟迟不输出感觉任务死掉了原因你设了 EventTime 窗口但数据源没有正确分配时间戳和 Watermark或者数据本身的 Kafka partition 数大于 Flink 并行度导致部分 source 子任务收不到数据Watermark 永远停在初始值。窗口的触发条件是 Watermark 超过窗口结束时间如果 Watermark 不涨窗口永远不触发。解决先用日志确认每并行源是否都有数据。然后在assignTimestampsAndWatermarks之后立刻打印 watermarks 看看。也可以临时改成 ProcessingTime 验证逻辑是否通但要记得改回来。一个重要技巧对于空闲 partition配置WatermarkStrategy.withIdleness(Duration.ofMinutes(1))让空闲 source 自动推进防止卡死。6.3 现象HBase 写入性能差任务背压上涨最终 OOM原因很多人用table.put()单条写入或者开着 HBase 的 autoflush 逐条刷写海量推荐结果涌入时 region server 压力巨大。另一个原因是 RowKey 设计用纯 userId 开头导致所有写落在同一个 region 上形成热点。解决改用 BufferedMutator 批量写把写缓冲调到 8MB 左右。RowKey 加盐MD5(userId).substring(0,2) _ userId让前缀分布到不同 region。如果项目里允许重启还可以预拆分 HBase 表的 Region按前缀指定 split keys比自动分裂稳定得多。6.4 现象Flink 状态无限增长磁盘被 RocksDB 占满原因你没有清理过期的 Keyed State。比如用ListState保存用户行为序列如果只 add 不清理或者定时器永远不触发状态就会无限膨胀。用户量大、窗口粒度细的时候尤其恐怖。解决对每个 state 明确声明 TTL。StateTtlConfig设置存活时间比如 30 分钟并配置 cleanUp 策略StateTtlConfig ttlConfig StateTtlConfig .newBuilder(Time.minutes(30)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .cleanupInRocksdbCompactFilter(1000) .build(); ListStateDescriptorUserBehavior descriptor new ListStateDescriptor(behavior-seq, UserBehavior.class); descriptor.enableTimeToLive(ttlConfig);逻辑说明这段配置让behavior-seq状态在 30 分钟内没有读写就过期并且开启 RocksDB 压缩清理避免删除逻辑只在过期时惰性变色、实际不回收磁盘。参数注意cleanupInRocksdbCompactFilter(1000)的 1000 表示每处理 1000 条就查一次过期键值太小影响性能太大则状态清理不及时100 到 1000 是常见区间。6.5 现象Spring Boot 应用查询推荐结果频繁超时原因这是系统衔接问题。推荐结果存在 HBaseSpring Boot 服务每次请求都去查 HBase但 HBase 连接没有复用每次新建 Connection光 RPC 建链就占了几百毫秒。并发一高就超时。解决在 Spring Boot 里使用 HBase 连接池缓存 Connection 和 Table 对象的复用。另外在查询层加一层缓存——用 Caffeine 本地缓存推荐结果 30 秒接入层用 Redis 缓存 1 分钟。虽然推荐希望实时但用户对 30 秒到 1 分钟的延迟是无感的缓存能极大降低 HBase 压力。Spring Boot 整合 Flink 的另一个点是不要在 Web 容器线程里直接启动 Flink 作业应该把 Flink 任务独立部署Spring Boot 只负责任务管理和结果查询。6.6 现象Flink 作业恢复后推荐结果出现重复或丢失原因检查点恢复的 exactly-once 语义在 Sink 层不一定保证。Flink 可以对状态做 exactly-once但写入 HBase 这个外部系统时如果没做幂等设计重放就会产生重复数据。HBase 的 Put 本身是覆盖写相同 RowKey 相同列会覆盖所以重复概率低但如果你把时间戳拼进 RowKey重放就会生成新版本导致重复推荐。解决RowKey 里不要拼时间戳或只拼分钟级时间戳或者把 userId 唯一化——同一个用户多次计算推荐结果写同一 RowKey用最新写覆盖旧值。这样重放后最终结果是一致的。再配合检查点配置设置合适的间隔比如 60 秒一次让恢复粒度可控。7. 进阶验证与工程化习惯把推荐结果从「能跑」调成「准确」验证一个实时推荐系统不能只看任务不报错得看推荐质量。一个我常用的验证方法是离线回放把data目录的历史日志按时间顺序灌进 Flink在一个时间点停住流记录当时的推荐结果然后和真实用户后续点击的商品对比计算 hit rate。这个流程不需要额外写很多代码只要在 Flink 里加一个侧输出side output把每次窗口计算的结果同时发到 Kafka 或者打印到日志里就能收集数据做评估。侧输出的做法很简单OutputTagRecommendResult recommendTag new OutputTag(recommend, Types.POJO(RecommendResult.class)); // 在窗口 process 里 ctx.output(recommendTag, result); // 在 main 流之外获取 DataStreamRecommendResult recommendSideStream mainStream.getSideOutput(recommendTag); recommendSideStream.addSink(new KafkaSink(...));逻辑说明OutputTag用于标记旁路流适合把「主业务数据」和「监控评估数据」分开。主流程继续往下游写 HBase旁路把推荐结果单独发到 Kafka 供离线评估消费。这种旁路设计在生产里非常实用不会干扰主链路。参数注意Tag 要用Types.POJO明确类型否则序列化比较麻烦侧输出流建议单独设定并行度不要和主流程抢资源。验证看两个指标最直接推荐的多样性即用户看到的商品不能总是同一个类目下的热销品覆盖率即推荐列表里有多少不同的商品。这个系统如果你能跑完你会发现窗口越大覆盖率越高但即时响应越差。此时就可以用参数试验把窗口从 10 分钟改成 5 分钟观察 hit rate 变化。我一般会在实际数据上对比三组窗口5、10、15 分钟每组跑一小时用旁路输出的日志计算指标取最优。哈希前缀的盐值也是关键RowKey 加盐长度会影响 HBase Region 分布如果发现某个 Region 数据量远大于平均试试改前缀长度。另外不要忽略 Flink UI 的监控面板。重点关注Watermark曲线它应该是一条稳步上升的锯齿线再看currentInputWatermark与 EventTime 的差值那代表了你系统的实际乱序延迟。如果差值和你的 Watermark 允许延迟接近说明数据乱序严重需要适当增大容忍时间或者优化采集端。而backPressure指标如果显示 high别急着加资源先查 HBase Sink 是不是瓶颈——我遇见太多次「背压上涨」其实是 HBase 批量写缓冲设置太小换成 16MB 后马上平缓。有一种执念我觉得新人最容易有总想把推荐算法换成复杂的深度模型。这份工程用的是可解释的相似度推荐这恰恰是生产环境里最稳的方案。模型可以慢慢迭代但实时链路只要可靠业务就赢了。从那以后我每次部署这套系统都会强制走一遍「回放评估 HBase 热点检查 状态 TTL 确认」这三步才上线。因为踩过状态膨胀导致磁盘爆掉的坑也见过 RowKey 热点让 HBase 整个集群卡死的夜晚这些细节用一次血泪换一次教训比任何炫技算法都值钱。希望帮到你。本文还有配套的精品资源点击获取