资讯动态

基于Flink的商品实时推荐系统:从离线T+1到秒级更新实战

发布时间:2026/10/9 1:13:11 来源:尧图企业网站定制
简介这是一套基于Flink的商品实时推荐系统完整项目资料面向计算机相关专业的在校学生、教师及企业开发人员可用于毕业设计、课程设计、项目立项演示或大数据技术进阶学习。资源包共47个文件约247KB以34个Scala源码文件为核心配合SQL脚本、properties配置、xml依赖文件、HBase建表语句、Kafka模拟数据及说明文档覆盖实时推荐系统的数据采集、流式计算与存储落地等关键环节。项目已通过测试运行功能完整并附有详细文档与README说明便于快速理解整体架构与模块划分。目前已有57人学习下载。读者可据此掌握Flink实时计算与推荐逻辑的实现思路参考HBase存储与Kafka数据模拟的配置方式也可在现有代码基础上修改扩展完成自己的毕设或课设任务是一份实用性较强的学习与开发参考。1. 基于 Flink 的商品实时推荐系统从离线 T1 到秒级更新的分水岭商品推荐这件事离线批处理时代大家都很熟每天凌晨跑一遍 Spark 任务把用户历史行为灌进 ALS 或者 ItemCF算出推荐结果写回 MySQL第二天早上生效。这套流程稳定、好维护但有个致命问题——用户上午十点刚点了一堆母婴用品系统到第二天才知道他可能是个新手爸爸。实时推荐系统要解决的就是这个时间差用户行为发生后几百毫秒到几秒内推荐列表就跟着变。Flink 之所以成为这类系统的首选计算引擎核心在于它原生支持事件时间、状态管理和精确一次语义能把「用户点击 → 特征更新 → 召回排序 → 结果下发」这条链路做成真正的流式管道而不是微批模拟。这套方案适合已经有离线推荐基础、想往实时方向演进的团队也适合刚接触 Flink 但需要一个完整落地场景来练手的工程师。接下来我会按「数据从哪来 → 特征怎么算 → 推荐怎么出 → 坑在哪」的顺序把这条链路拆开讲清楚。2. 商品实时推荐的数据链路从埋点到 Flink 的完整拓扑2.1 为什么不能直接把离线推荐那套搬到 Flink 上很多团队第一次做实时推荐最容易犯的错是把离线那套特征工程原封不动搬到 Flink 里跑。离线场景下你可以全量扫描用户过去 30 天的行为表随便 join、随便 group by反正跑完就完事。但流式场景下数据是无界且持续到达的你没法「等所有数据到齐再算」。这就意味着两件事必须重新设计一是特征的窗口定义二是状态的存储和过期策略。常见做法是把用户行为分成短期兴趣和长期兴趣两层。短期兴趣用滑动窗口比如最近 10 分钟内的点击、加购、收藏行为窗口每 30 秒滑动一次输出用户对商品类目的实时偏好分数。长期兴趣则用 Flink 的状态后端维护一个用户画像向量每天从离线数仓同步一次基线实时行为只做增量更新。这样既保证了实时性又不会因为纯流式计算导致特征抖动过大。另一个关键选型是消息队列。Kafka 几乎是标配但要注意分区策略如果按 userId 分区同一用户的行为会落到同一分区方便做 keyed state但如果某个头部主播带货导致单用户行为暴涨会出现数据倾斜。我一般会先用 userId 的 hash 做分区同时在 Flink 侧加一层本地预聚合把同一用户 100ms 内的多次点击合并成一次处理减轻下游压力。2.2 用 Flink DataStream API 搭一条最小可跑的行为采集管道下面这段代码是一个最小化的行为采集管道从 Kafka 读取用户行为 JSON解析后按 userId 分组用 10 分钟滑动窗口统计每个用户对每个类目的点击次数结果输出到另一个 Kafka topic 供下游召回使用。# 注意这是 PyFlink 的写法Java/Scala 版本逻辑一致但 API 略有不同 from pyflink.datastream import StreamExecutionEnvironment, TimeCharacteristic from pyflink.datastream.window import SlidingEventTimeWindows, Time from pyflink.common import WatermarkStrategy, Duration from pyflink.datastream.connectors import FlinkKafkaConsumer, FlinkKafkaProducer from pyflink.common.serialization import SimpleStringSchema import json env StreamExecutionEnvironment.get_execution_environment() env.set_stream_time_characteristic(TimeCharacteristic.EventTime) env.enable_checkpointing(60000) # 60秒一次 checkpoint保证精确一次 # 1. 从 Kafka 读取原始行为日志 kafka_source FlinkKafkaConsumer( topicsuser_behavior_raw, deserialization_schemaSimpleStringSchema(), properties{ bootstrap.servers: kafka-broker:9092, group.id: realtime_rec_group } ) # 2. 解析 JSON 并分配事件时间戳用行为发生时间不是处理时间 def parse_and_assign_watermark(json_str): data json.loads(json_str) # 关键用 event_time 字段作为事件时间允许 5 秒乱序 return data stream env.add_source(kafka_source) \ .map(parse_and_assign_watermark) \ .assign_timestamps_and_watermarks( WatermarkStrategy.for_bounded_out_of_orderness(Duration.of_seconds(5)) .with_timestamp_assigner(lambda x, _: x[event_time]) ) # 3. 按 userId categoryId 分组做 10 分钟滑动窗口每 30 秒滑一次 windowed stream \ .key_by(lambda x: (x[user_id], x[category_id])) \ .window(SlidingEventTimeWindows.of(Time.minutes(10), Time.seconds(30))) \ .aggregate(ClickCountAggregator()) # 自定义聚合函数统计点击次数 # 4. 结果写回 Kafka供召回层消费 kafka_sink FlinkKafkaProducer( topicuser_category_click_10min, serialization_schemaSimpleStringSchema(), producer_config{bootstrap.servers: kafka-broker:9092} ) windowed.map(lambda r: json.dumps(r)).add_sink(kafka_sink) env.execute(Realtime User Behavior Pipeline)这段代码里有几个参数需要根据实际数据量调整。enable_checkpointing(60000)里的 60000 是毫秒表示每分钟做一次状态快照。如果状态很大比如用户画像向量维度很高可以调到 120000 甚至 300000但代价是故障恢复时重放的数据更多。for_bounded_out_of_orderness(Duration.of_seconds(5))里的 5 秒是乱序容忍度移动端网络延迟通常在这个范围内如果发现窗口结果经常漏数据可以适当加大到 10 秒但窗口输出会相应延迟。SlidingEventTimeWindows.of(Time.minutes(10), Time.seconds(30))表示窗口长度 10 分钟、滑动步长 30 秒。这意味着每个用户每 30 秒就会输出一次过去 10 分钟的点击统计。步长越小实时性越好但计算和下游写入压力越大。我一般会先设 30 秒观察 Kafka 下游消费延迟如果消费跟不上就调到 60 秒。2.3 状态后端选型RocksDB 还是 HashMapFlink 的状态后端直接决定了你能维护多大的用户特征。HashMapStateBackend 把状态放在 JVM 堆内存里读写快但受限于 TaskManager 的内存上限通常只能撑几百万个 key。RocksDBStateBackend 把状态存在本地磁盘支持增量 checkpoint能撑几十亿 key但读写有序列化开销。对于商品实时推荐如果用户量在百万级以内且特征维度不高比如每个用户只存最近 50 个交互商品HashMapStateBackend 完全够用延迟更低。但如果要做全量用户画像的实时更新或者用户量上千万就必须上 RocksDB。配置方式是在flink-conf.yaml里设置state.backend: rocksdb state.backend.incremental: true state.checkpoints.dir: hdfs://namenode:8020/flink/checkpoints state.backend.rocksdb.memory.managed: truestate.backend.incremental: true是 RocksDB 的增量 checkpoint只上传变化的 SST 文件能把 checkpoint 时间从分钟级降到秒级。state.backend.rocksdb.memory.managed: true让 RocksDB 的内存由 Flink 统一管理避免和算子内存打架。这两个参数是我踩过坑之后必开的。3. 实时特征计算与召回排序把用户行为变成推荐结果3.1 实时特征的两层结构短期兴趣和长期画像实时推荐的特征不能只靠流式计算也不能只靠离线数仓。我的做法是分两层第一层是 Flink 实时算出来的短期兴趣分第二层是从 HBase 或 Redis 读取的长期画像向量。两者在召回层做加权融合。短期兴趣分的计算逻辑是用户对某个类目的点击、加购、收藏、下单行为分别赋予不同权重比如点击 1 分、加购 3 分、收藏 2 分、下单 5 分在 10 分钟滑动窗口内累加再除以该用户的总行为次数做归一化。这样得到的分数在 0 到 1 之间代表用户当前对这个类目的偏好强度。长期画像向量则是离线任务每天更新一次存在 HBase 里rowkey 是 userId列族里存各个类目的历史偏好分。Flink 在召回阶段通过异步 IO 查询 HBase把长期向量和短期分数做加权求和。权重通常短期占 0.6、长期占 0.4但大促期间会调成 0.8 比 0.2因为大促时用户兴趣变化极快历史画像参考价值下降。3.2 用 Flink SQL 做实时召回JDBC 连接器查 MySQL 商品池召回阶段需要从商品池里筛选候选集。商品池通常存在 MySQL 里包含商品 ID、类目、价格、库存等字段。Flink SQL 的 JDBC 连接器可以直接把 MySQL 表映射成一张维表和实时行为流做 join。下面是一个完整的 Flink SQL 示例-- 1. 创建 Kafka 行为源表 CREATE TABLE user_behavior ( user_id BIGINT, item_id BIGINT, category_id BIGINT, behavior_type VARCHAR, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic user_behavior_raw, properties.bootstrap.servers kafka-broker:9092, properties.group.id realtime_rec_sql, format json, scan.startup.mode latest-offset ); -- 2. 创建 MySQL 商品维表注意这是 JDBC 连接器不是 CDC CREATE TABLE item_dim ( item_id BIGINT, category_id BIGINT, item_name VARCHAR, price DECIMAL(10,2), stock INT, PRIMARY KEY (item_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://mysql-host:3306/recommend, table-name item_pool, username flink_reader, password ******, lookup.cache.max-rows 10000, lookup.cache.ttl 10min ); -- 3. 实时 join行为流关联商品维表过滤掉无库存商品 CREATE TABLE recall_result AS SELECT b.user_id, b.item_id, i.item_name, i.category_id, i.price, b.event_time FROM user_behavior b JOIN item_dim FOR SYSTEM_TIME AS OF b.event_time AS i ON b.item_id i.item_id WHERE i.stock 0;这里有几个参数直接决定系统能不能稳定跑。lookup.cache.max-rows 10000表示维表缓存最多存 1 万行商品数据如果商品池有几十万 SKU这个值要调大但会占用更多 TaskManager 内存。lookup.cache.ttl 10min是缓存过期时间商品价格和库存变化不频繁的话可以设长一点减少对 MySQL 的查询压力。FOR SYSTEM_TIME AS OF b.event_time是 Flink SQL 的时态表 join 语法它保证用行为发生时刻的商品快照来关联而不是用当前最新值。这个细节在价格敏感场景下很重要——用户点击时看到的价格和推荐时展示的价格必须一致否则会出现「点击时 99 元推荐列表里变 129 元」的翻车情况。3.3 排序层用 Flink 做实时打分还是交给外部服务召回之后需要排序。排序模型通常是一个深度学习模型比如 DIN、DeepFM推理需要 GPU 或者高性能 CPU。Flink 本身不适合做模型推理所以常见做法是 Flink 只负责召回和粗排把候选集写入 Redis 或 Feature Store再由外部排序服务拉取候选集和特征做精排。但有一种情况可以在 Flink 里做轻量级排序当排序规则是简单的线性加权比如短期兴趣分 × 0.6 商品热度分 × 0.3 价格偏好分 × 0.1完全可以用 Flink SQL 的 UDF 实现。这样省去了外部服务调用端到端延迟能控制在 200ms 以内。我一般会先用 Flink SQL 做一版规则排序上线验证链路通畅后再把精排模型接进来替换。3.4 把结果写回 MySQL 供前端查询JDBC Sink 的批量参数推荐结果最终要落到 MySQL 或者 Redis 供前端接口查询。如果用 Flink SQL 写 MySQLJDBC Sink 的批量参数必须调否则每条结果都单独 insertMySQL 会被打爆。CREATE TABLE rec_result_sink ( user_id BIGINT, item_id BIGINT, score DOUBLE, rec_time TIMESTAMP(3), PRIMARY KEY (user_id, item_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://mysql-host:3306/recommend, table-name realtime_rec_result, username flink_writer, password ******, sink.buffer-flush.max-rows 500, sink.buffer-flush.interval 2s, sink.max-retries 3 );sink.buffer-flush.max-rows 500表示攒够 500 条才批量写入sink.buffer-flush.interval 2s表示最多等 2 秒也会强制写入。这两个参数配合使用既保证吞吐又控制延迟。sink.max-retries 3是写入失败重试次数MySQL 主从切换时能扛一下短暂不可用。4. 避坑与排查实时推荐系统上线后最容易翻车的五个地方4.1 现象Flink 任务频繁重启日志报 JDBC 连接超时原因JDBC 连接器在维表 join 时会为每个并发子任务创建连接池如果并发度高比如 100 个并行度MySQL 的 max_connections 直接被占满后续连接请求全部超时。解决在 JDBC 连接器参数里加connection.max-retry-timeout和connection.pool.size把每个子任务的连接池大小限制在 2 到 3同时和 DBA 确认 MySQL 的最大连接数至少是 Flink 并行度的 3 倍。更稳妥的做法是维表数据量不大时直接广播到 Flink 状态里不走 JDBC 查询。4.2 现象窗口计算结果比预期少部分用户行为丢失原因水位线设置过于激进乱序数据到达时窗口已经关闭迟到数据被丢弃。移动端网络抖动、Kafka 分区重平衡都会导致乱序。解决把for_bounded_out_of_orderness从 5 秒调到 10 秒同时在窗口上允许迟到数据window(...).allowed_lateness(Time.minutes(1))并配置侧输出流收集迟到超过 1 分钟的数据离线补跑。注意 allowed_lateness 会增加状态保留时间RocksDB 的磁盘占用会上升。4.3 现象Checkpoint 持续失败报 RocksDB 写入超时原因RocksDB 的本地磁盘 IO 被打满通常是状态太大或者磁盘用的是网络存储如 NFS写入延迟高。解决把state.backend.rocksdb.localdir指向本地 SSD 盘不要用 NFS。同时开启state.backend.incremental: true减少每次 checkpoint 的数据量。如果状态确实太大考虑把部分历史特征从 Flink 状态里挪到外部存储HBase/RedisFlink 只维护最近窗口的状态。4.4 现象推荐结果里出现已下架商品原因JDBC 维表缓存 TTL 设得太长商品下架后缓存还没过期Flink 仍然用旧数据做 join。解决把lookup.cache.ttl从 10 分钟降到 1 分钟或者改用 MySQL CDC 连接器实时捕获商品表变更把维表做成动态更新的。CDC 方式延迟更低但需要 MySQL 开启 binlog 并配置 Debezium。4.5 现象大促期间 Kafka 消费延迟飙升推荐结果滞后十几分钟原因行为数据量突增Flink 算子并行度不够或者 Kafka 分区数太少导致消费能力受限。解决提前把 Kafka topic 分区数扩到 Flink 最大并行度的 2 倍Flink 任务开启自动扩缩容如果用的是云托管版本或者手动调大并行度。另外可以在 Flink 上游加一层预聚合把同一用户 100ms 内的多次点击合并减少下游处理量。大促前做一次全链路压测确认瓶颈在哪个环节。5. 进阶技巧用 Flink CDC 把 MySQL 商品变更实时同步到推荐特征库前面讲的 JDBC 维表 join 有个天然缺陷它是拉取模式商品变更后要等缓存过期才能生效。如果商品价格频繁调整或者库存变化需要秒级反映到推荐结果里就得用 Flink CDC 把 MySQL 的 binlog 实时捕获成流和用户行为流做双流 join。具体做法是用mysql-cdc连接器创建商品变更源表按 item_id 分组用Temporal Table Join或者Interval Join和行为流关联。这样商品价格一变推荐结果里的价格字段在秒级内跟着变。配置上需要注意 MySQL 的 binlog 格式必须是 ROW且binlog_row_image设为 FULL否则 CDC 拿不到完整的变更前后字段。验证 CDC 链路是否正常我一般会写一个简单的对账任务每隔 5 分钟统计一次 MySQL 商品表和 Flink 侧商品状态表的记录数如果差异超过阈值就告警。这个对账逻辑用 Flink SQL 的TUMBLE窗口加COUNT就能实现不需要额外写代码。还有一个容易忽略的点是 CDC 的初始快照阶段。如果商品表有上千万行全量快照会跑很久期间 binlog 还在持续产生。Flink CDC 支持增量快照incremental snapshot把全量阶段分成多个 chunk 并行读取能显著缩短启动时间。开启方式是设置scan.incremental.snapshot.enabled true同时确保 MySQL 用户有RELOAD和REPLICATION SLAVE权限。我自己在第一次做实时推荐时最大的教训是低估了状态管理的复杂度。当时觉得 Flink 的 keyed state 用起来很简单结果上线后状态膨胀到几百 GBcheckpoint 一次要十几分钟任务稳定性极差。后来把状态拆成「热状态」和「冷状态」热状态只保留最近 1 小时的数据在 Flink 里冷状态全部下沉到 HBasecheckpoint 时间直接降到 30 秒以内。这个习惯我一直保持到现在任何要进 Flink 状态的数据先问一句「它真的需要留在状态里吗」。希望帮到你。本文还有配套的精品资源点击获取

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑