
简介这是一套基于 Flink 实现的商品实时推荐系统项目源码面向大数据实时计算、推荐系统入门及课程设计人群可解决热度统计、用户画像、实时推荐请求响应等典型问题。系统利用 Flink 统计商品热度并写入 Redis 缓存分析日志后将画像标签与实时记录存入 HBase推荐时按用户画像重排热度榜再叠加协同过滤和标签推荐为每个商品补充关联产品最终生成个性化列表。资源共 109 个文件以 68 个 Java 源码为主同时包含 SQL 建表语句、HBase 建表脚本、Kafka 模拟数据脚本、YAML/Properties 配置和 HTML 前端页面压缩包仅 3.74MB整体结构清晰。已有 130 人学习适合参考核心实现类如 RecommandServiceImpl梳理工程结构也可作为二次开发或实战演练的起点。1. 实时推荐系统为什么绕不开 Flink一次缓存雪崩换来的认知做推荐的人最容易犯的错是把精力全压在离线模型和算法调参上结果模型还没上线数据管道先崩了。我接手这套基于 Flink 实现的商品实时推荐系统时线上正卡在一个诡异的故障里热度榜三个小时没更新推荐接口却照常返回因为读的是 Redis 里的旧缓存。后来调了两天才发现问题出在上游统计作业的背压跟算法一点关系没有。这套系统的链路其实很清晰Flink 消费点击日志实时统计商品热度写入 Redis 缓存日志分析的结果落到 HBase存用户画像和实时行为记录用户发起推荐请求后系统按画像重排序热度榜再用协同过滤和标签推荐模块给榜单商品补关联产品最后返回个性化列表。它适合两类人想找完整可复现项目练手的 Flink 学习者和技术选型时不确定链路怎么搭的从业者。下面按真实实现来讲每段代码都能直接搬。2. Flink 全链路架构从点击日志到热度榜的一条数据管道2.1 链路选型为什么是 Flink Redis HBase先交代选型因为这三个组件的组合不是拍脑袋定的。流处理框架对比过 Spark Streaming 和 FlinkSpark Streaming 是微批模型一批数据攒够才处理批间隔再小也是批延迟做不到秒级Flink 是事件驱动每条日志进来就能触发计算窗口只用来圈定统计范围。商品热度这种指标榜单更新的快慢直接决定推荐质量所以我选了 Flink。Redis 的角色是热度榜缓存。热度榜本质是一个按分数排序的集合Redis 的 ZSet 用 O(logN) 就能完成写入和 TopN 查询一条 ZREVRANGE 命令就能取出前 100读操作落在内存里QPS 再高也不慌。换成一个关系型数据库每次请求都走 ORDER BY热榜接口一上线就会被查询打崩。HBase 存的是用户画像和实时行为记录。这部分数据特点是量大、维度杂、更新频繁推荐服务要按 userId 精准读取。HBase 的列族设计让画像属性、标签权重、行为明细分开存放行键查询速度快也不会占用 Redis 的内存。三个组件各管一段Flink 算、Redis 存热榜、HBase 存画像职责不重叠。2.2 热度统计作业Kafka Flink 窗口聚合 Redis Sink热度统计的原理用 Flink 官方那个“词频统计初体验”来类比最好懂单词变成商品 ID出现次数变成加权热度输出从控制台换成了 Redis。下面是我在线上跑的作业骨架去掉监控和告警后的核心逻辑。输入是 Kafka 的 click_log 主题每条日志包含 userId、productId、action、time 四个字段。// HotProductJob.java import java.time.Duration; import java.util.Properties; public class HotProductJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(4); // 每 60 秒做一次 checkpoint状态量小默认配置足够 env.enableCheckpointing(60_000); Properties kafkaProps new Properties(); kafkaProps.setProperty(bootstrap.servers, kafka-1:9092,kafka-2:9092); kafkaProps.setProperty(group.id, hot-product-calc); // 重启后从最早位点补算防止热度榜出现空白 kafkaProps.setProperty(auto.offset.reset, earliest); DataStreamProductClick clicks env .addSource(new FlinkKafkaConsumer(click_log, new SimpleStringSchema(), kafkaProps)) .map(line - ProductClick.parse(line)) .assignTimestampsAndWatermarks( // 允许乱序延迟 5 秒超过的只能等下轮窗口 WatermarkStrategy.ProductClickforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((click, ts) - click.getTs()) ); DataStreamProductScore hot clicks .keyBy(ProductClick::getProductId) .window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1))) .aggregate(new CountAgg(), new WindowScore()); FlinkJedisPoolConfig redisConf new FlinkJedisPoolConfig.Builder() .setHost(redis-1) .setPort(6379) .setMaxTotal(16) .setMaxIdle(8) .build(); hot.addSink(new RedisSink(redisConf, new HotScoreMapper())); env.execute(hot-product-job); } }这段代码里几个参数要盯住。SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1))表示窗口长度 10 分钟、每 1 分钟滑动一次也就是说热度榜每分钟刷新但统计的是过去 10 分钟的累积值。forBoundedOutOfOrderness(Duration.ofSeconds(5))是最多容忍 5 秒乱序如果 Kafka 分区数多这个值要按线上日志到达延迟调一般先给 3 到 5 秒压测后看延迟数据再收紧。CountAgg 负责把同窗口内同一商品的行为加权累加。行为类型和权重的映射并不在作业里写死而是从日志里的 action 字段取点击 1、加购 3、下单 5、分享 4。这样加购和下单对热度的影响就能拉开差距不会出现“一个下单不如十个点击”的失真情况。public static class CountAgg implements AggregateFunctionProductClick, Long, Long { Override public Long createAccumulator() { return 0L; } Override public Long add(ProductClick click, Long acc) { return acc click.getWeight(); } Override public Long getResult(Long acc) { return acc; } Override public Long merge(Long a, Long b) { return a b; } }最后把结果写进 Redis用的是 RedisSink 加自定义 Mapper。getKeyFromData对应 ZADD 的 membergetValueFromData对应 score构造参数里的hot:products是 ZSet 的 key。public static class HotScoreMapper implements RedisMapperProductScore { Override public RedisCommandDescription getCommandDescription() { return new RedisCommandDescription(RedisCommand.ZADD, hot:products); } Override public String getKeyFromData(ProductScore score) { return String.valueOf(score.getProductId()); } Override public String getValueFromData(ProductScore score) { return String.valueOf(score.getScore()); } }部署这块提一句Flink 安装配置到部署不同版本差异很大我习惯用 Flink 1.15 以上的版本提交用 yarn-per-job 模式资源和任务隔离都比较干净。本地调试可以直接在 IDE 里跑env.execute()不需要先装一整套集群。热度统计本身没有状态恢复的特殊要求checkpoint 开 60 秒足以覆盖大部分故障场景。2.3 HBase 建表与画像存储从建表语句到 RowKey 设计日志分析后的画像标签和实时行为记录都落到 HBase。根目录下那份“hbase 建表语句”是整套系统的存储设计两个表解决三类数据建表脚本如下。# hbase-shell 建表脚本 create user_profile, {NAME base, VERSIONS 1, TTL 2592000, COMPRESSION SNAPPY}, {NAME tags, VERSIONS 1, TTL 2592000, COMPRESSION SNAPPY} create user_recent_behavior, {NAME recent, VERSIONS 1, TTL 1209600, COMPRESSION SNAPPY}表名列族存什么TTLuser_profilebase用户基础信息如注册渠道、年龄段、会员等级30 天user_profiletags用户兴趣标签及对应权重30 天user_recent_behaviorrecent最近 14 天的行为明细14 天TTL 是这里最关键的参数。画像标签虽然是长期累积的但用户兴趣变化快30 天不活跃的标签留着只会干扰重排序。标签权重在业务层做衰减存储层只保留新分数所以 VERSIONS 保持 1不开多版本。RowKey 直接用了 userId。表面看会有热点用户挤在同一个 Region但推荐场景是读多写少热点用户频繁读取反而让 Region 缓存命中率更高。等 userId 规模到亿级再考虑加盐现阶段用原生 userId 做行键就够了。画像标签的写入逻辑在 UserScoreServiceImpl 里。每次用户行为进来按规则算出新的标签分数然后以 put 方式写回 tags 列族。// UserScoreServiceImpl 写标签的核心片段 Put put new Put(Bytes.toBytes(userId)); put.addColumn(Bytes.toBytes(tags), Bytes.toBytes(category: categoryId), Bytes.toBytes(String.valueOf(newScore))); userProfileTable.put(put);列名用category:类目ID这种带前缀的拼接方便查询时按前缀过滤也能在一个列族里放下多种类型的标签而不用为每个标签类型建单独的列族。3. 画像重排序与关联商品把热度榜变成千人千面列表3.1 用户画像构建行为日志如何变成 HBase 标签热度榜是全局的所有用户看到的是同一个排序。要让每个用户看到不同的顺序就需要画像参与重排。画像不是静态属性而是动态兴趣标签用户最近 10 分钟高频点击的类目、最近 3 天反复浏览的品牌、加购但没下单的商品都会影响标签权重。标签的计算策略是“有衰减的加权累加”。每次行为事件进来按行为类型取权重再和旧分数加权合并公式是newScore oldScore × 0.8 weight。0.8 是衰减系数越大历史影响越重越小越偏向当前行为。实时场景下我一般给 0.8这样用户三天前的行为还能留一点痕迹但不会盖过今天的新动作。日志里每一条行为都要带时间戳和来源渠道否则无法做时间衰减也无法区分 Web 端和 App 端的行为差异。这个细节在整个链路里容易被忽略但直接影响画像质量后面重排序的准头全靠它。画像标签分两类一类是类目偏好标签前缀是category:一类是品牌偏好标签前缀是brand:。两类标签都存同一个 tags 列族靠列名前缀区分。推荐服务重排序时只关心类目标签因为热度榜上的商品必然有类目属性匹配类别直接品牌标签则留给后续的精细排序用。3.2 热度榜重排序UserScoreServiceImpl 的加权评分实现用户发起推荐请求后推荐服务先从 Redis 取热度榜 TopN再调 UserScoreServiceImpl 按画像重排。核心公式是finalScore heatScore × 0.6 profileScore × 0.4。热度分是全局的画像分是用户自己的两者加权后才构成这个用户视角下的商品排序。// UserScoreServiceImpl.rerank 核心逻辑 public ListProduct rerank(String userId, ListProduct hotList) { MapString, Double userTags hbaseDao.getUserTags(userId); double maxHeat hotList.stream() .mapToDouble(Product::getHeatScore).max().orElse(1.0); return hotList.stream().map(p - { // 归一化热度避免榜单尺度差异 double normHeat p.getHeatScore() / maxHeat; // 从画像里取该类目偏好分 Double tagScore userTags.get(category: p.getCategoryId()); double normTag tagScore null ? 0.0 : tagScore; p.setRankScore(normHeat * 0.6 normTag * 0.4); return p; }).sorted(Comparator.comparingDouble(Product::getRankScore).reversed()) .collect(Collectors.toList()); }参数取值说明heatWeight0.6全局热度权重热度可信度越高越偏向它tagWeight0.4用户画像权重冷启动时可以调低窗口长度10 分钟热度统计的滑动窗口重排规模前 200 个只重排榜单头部的 200 个商品两个细节值得说。一是热度归一化用的分母是当前榜单的 maxHeat而不是历史峰值否则冷启动期间所有商品分数都接近 0排序不稳定。二是 tagScore 为 null 时不能粗暴给 0因为新用户本来就没有标签给 0 会让这批商品全部靠在热度分上相当于没重排但冷启动用户本来就应该看热门这个行为其实是正确的兜底。3.3 协同过滤 标签推荐RecommandServiceImpl 的关联融合重排序只解决了“榜单顺序”每个商品本身还缺“关联商品”。这一步在 RecommandServiceImpl 里完成协同过滤和标签推荐两个模块各出一组关联品再按规则融合。// RecommandServiceImpl.buildResult 核心逻辑 public ListProduct buildResult(String userId, ListProduct rankList) { ListProduct result new ArrayList(); for (Product p : rankList) { result.add(p); // 协同过滤基于商品共现关系推荐同场景关联品 ListProduct cfItems cfModule.getRelated(p.getProductId(), 3); // 标签推荐基于用户画像和商品标签匹配来补位 ListProduct tagItems tagModule.getRelated(userId, p.getCategoryId(), 2); // 融合cf 优先tag 补位合计最多 3 个 result.addAll(merge(cfItems, tagItems, 3)); } return result; }协同过滤模块我用的是基于商品的 Item-CF核心是商品与商品的共现矩阵。共现矩阵离线算好放进 Redis线上按 productId 直接查不参与实时计算。标签推荐模块则把当前商品的类目标签和用户画像做匹配类目一致且用户偏好高的商品打分高。融合策略上CF 结果放在前面因为共现关系比标签推荐更直接标签结果只在 CF 结果不足时补位。每个主商品最多补 3 个关联品防止返回列表膨胀导致前端加载变慢。最终 buildResult 返回的就是新的商品列表前 20 个作为主榜单其余作为后续的流式补位内容。4. 避坑指南五个翻车现场与排查方法实时推荐系统能跑起来只是开始下面五个坑我都是实际踩过的按“现象→原因→解决”列出来直接对应到代码和参数。4.1 Flink JDBC 连接器异常连接池耗尽与参数修正现象作业做维表关联时运行几小时后开始报Connection is not available, request timed out重启作业能恢复几小时后又复现。原因Flink JDBC 连接器默认连接池很小我当时没设置池大小默认上限只有 8而作业并行度是 16高峰期所有 slot 同时访问维表连接池瞬间被打满后续请求一直等待到超时。解决连接池上限改成“并行度 × 2”起步我线上设了 32同时开启连接池的保活探测避免数据库侧把空闲连接回收后连接池还认为它可用。如果维表数据量能塞进内存更彻底的办法是用 Flink 的广播状态把全量维表加载进来从根上消除 JDBC 瓶颈。4.2 Sink Hive 表数据不入表分区提交与压缩参数现象Flink SQL 作业写 Hive 表TaskManager 日志显示任务正常结束数据文件也写了但去 Hive 查询时一条记录都查不到。原因Flink 写 Hive 的默认行为是写进 HDFS 临时目录只有触发分区提交后文件才对查询可见。我当时没设置sink.partition-commit.trigger默认按文件滚动提交而任务的数据量没攒够触发阈值导致数据一直卡在临时目录。解决给表配置分区提交策略用 partition-time 触发CREATE TABLE user_tag_hive ( user_id STRING, tag STRING, score DOUBLE ) PARTITIONED BY (dt STRING) WITH ( connector hive, sink.partition-commit.trigger partition-time, sink.partition-commit.delay 1 min, sink.partition-commit.watermark-time-zone Asia/Shanghai );partition-commit.delay给 1 分钟是留给数据从临时目录移到正式分区的缓冲时间给 0 会导致文件刚写完就被提交读端可能看到不完整的数据。4.3 自定义 Data Source并行度设置与状态恢复现象自定义 Source 读取文件时并行度设 1 一切正常改成 4 后出现重复消费和漏数据开启 checkpoint 恢复后还会丢一段数据。原因我最初继承的是 SourceFunction没实现 CheckpointedFunctionFlink 不知道每个 Source 实例消费到哪里。并行度变大后多个实例各读各的没有协调机制offset 也不统一。解决改继承 RichParallelSourceFunction在 snapshotState 里保存每个实例的消费位点restoreState 时从位点恢复。如果 Source 本身就是 Kafka直接用 Flink 内置的 FlinkKafkaConsumer它的 offset 管理已经处理好了这些细节不需要自己写 Source。自写 Source 只适合接入文件、数据库增量等 Kafka 覆盖不到的场景。4.4 Redis 缓存穿透热榜 key 的失效风暴现象热榜接口 P99 延迟从 5ms 飙到 300ms 以上Redis 的 CPU 接近 100%同时有部分请求穿透到了 HBase。原因热度榜只有一个固定 keyhot:products每轮窗口的写入全部压在同一个 key 上Redis 单线程模型下所有命令排队读请求被阻塞。TTL 到点后新数据还没生成读请求从 Redis 拿不到值就穿透到 HBase把 HBase 流量也带起来了。解决把热榜拆成多个分片 key比如hot:products:shard:0到hot:products:shard:9写入时均匀分发读取时并发取多个分片再合并排序TTL 用固定值加随机偏移避免同一时刻集体失效榜单重建前先预热新 key预热完成再切换读路径旧 key 到点自然过期。4.5 Flink CDC Pipeline 部署binlog 格式与连接器版本现象用 Flink CDC Pipeline 从 MySQL 同步用户表作业提交后报TableNotFoundException还有一次报 binlog 解析相关异常。原因两个问题叠加。一是 MySQL 的 binlog_format 不是 ROW 模式CDC 只能消费 ROW 格式的 binlog二是 flink-cdc 连接器版本和 Flink 运行时版本不匹配导致作业提交时类加载失败。解决先确认 MySQL 的binlog_formatROW和binlog_row_imageFULL这两项是 CDC 的前提然后核对 flink-cdc 连接器版本与 Flink 版本的对应关系Flink 1.15 对应 cdc 2.3.x 这一档按官方版本映射表对齐不要凭感觉选新版本。另外同步的表必须有主键CDC 的 snapshot 阶段需要按主键分批读取。5. 进阶用火焰图排查背压把推荐时延压到秒级以内线上排查背压我最常用的手段是抓 TaskManager 的火焰图。Flink Web UI 能告诉你哪个算子是瓶颈但看不到具体是哪段代码拖慢了线程火焰图能直接让热点栈现形。先用 jps 找到 TaskManager 进程号再采样。采样间隔 5ms持续 60 秒基本能覆盖一个完整的窗口周期。jps -l | grep TaskManager ./profiler.sh -d 60 -o flamegraph -i 5ms pid生成的火焰图是一个网页最宽的栈就是热点。热度统计作业如果最宽栈落在 Netty 的读写上说明流量阻塞在 Kafka 消费端优先扩并行度而不是调窗口如果栈落在 Redis 客户端等待上去查连接池参数。5.1 火焰图定位瓶颈排查对象火焰图热点栈常见根因对应调整热度统计作业Netty IO 读写Kafka 消费背压扩并行度、调大 fetch.max.bytes推荐服务作业HBaseScanner 扫描单行数据过大拆分列族、只读需要的列重排序逻辑Collectors / 排序候选列表过大只取前 200 个参与重排5.2 时延打点验证优化效果火焰图定位瓶颈后还需要一个能对拍的指标验证优化效果。我在推荐服务入口加了一行耗时打点把 buildResult 的完整耗时打到日志里。long start System.currentTimeMillis(); ListProduct result recommendService.buildResult(userId, hotList); long cost System.currentTimeMillis() - start; log.info(recommend_cost{}ms, items{}, cost, result.size());最典型的一次推荐服务 P99 延迟 2 秒火焰图最宽的栈是 HBaseScanner排查后发现 tags 列族里一个用户的行存了 2000 多个标签一次扫描要读几 MB 数据。把标签按类目拆到独立列只读重排需要的类目列延迟从 2 秒降到 300ms 以内。从那以后我每次调这类实时推荐的延迟问题都强制先抓一把火焰图再动参数不再凭感觉扩并行度。整个项目的源码我整理成了资源包HBase 建表语句、UserScoreServiceImpl、RecommandServiceImpl 和页面示例都在里面改掉 Kafka、Redis、HBase 三处连接地址就能复现这条链路。希望帮到你。本文还有配套的精品资源点击获取