
简介基于Flink实现的商品实时推荐系统源码包面向大数据流计算与推荐系统开发者完整演示了从日志采集、热度统计到个性化推荐的生产级实现。系统依托Flink实时计算商品热度并写入Redis缓存同时将用户画像与行为记录存入HBase通过协同过滤与标签双重模块为用户重排序榜单并补充关联商品可应用于电商个性化推荐、实时营销等场景。这套源码包含完整项目结构与可运行代码核心逻辑涵盖推荐服务、用户评分、商品热度统计等环节。资源共109个文件以68个Java源码为主辅以XML配置、SQL建表语句、Properties和YAML配置文件以及HTML展示页面、HBase建表语句、Kafka模拟数据脚本等压缩包大小3.74MB目录结构清晰便于按模块查阅。已有129人学习下载。对希望掌握Flink、Redis、HBase整合以及推荐算法落地的读者而言这套项目既能用于代码参考也可作为二次开发的起点。1. 实时推荐为什么要选Flink下午浏览过一双越野跑鞋晚上再打开App推荐流要能接住这个信号要么推同款要么推同类相关商品。要做到“接住”只靠凌晨跑一次的离线调度不够用点击、加购、下单这类信号的价值随时间快速衰减。用Flink做商品实时推荐本质是把用户行为流实时转换成特征、相似度和候选集在秒级更新召回结果。这里要处理的不只是推荐算法本身还有行为数据接入、窗口特征计算、结果写入和链路排错。这套方案适合已经用 Hive 或 Spark 做过离线推荐、正在规划实时链路的团队也适合想弄清 Flink 在推荐场景里到底怎么落地的工程师。2. Flink在实时推荐中的计算模型与数据基础实时推荐不是把离线管道跑快一点。离线推荐拿到的是 T1 全量数据算的是用户过去一个月偏好总账实时推荐关心的是“最近一小时”甚至“最后三次点击”。两者不是替代关系而是互补离线模型给基础排序实时特征做增量偏置。2.1 实时推荐与离线推荐的本质差别先看清楚差别才知道每条数据该往哪条管道放。维度离线推荐Flink 实时推荐数据范围全量历史行为最近几分钟到几小时特征更新小时级或天级秒级到分钟级算法目标训练模型、批量计算相似度快速反映短期意图结果时延分钟到小时秒级故障恢复重跑任务基于 checkpoint 增量恢复离线推荐算的是“稳定偏好”比如用户过去一个月常买咖啡豆可以做个长期画像。但用户半小时前刚搜索过帐篷这个信号如果等到明天才进特征转化机会就没了。实时推荐要捕捉的就是这种“短期意图急转弯”。常见做法是双通道设计离线结果做主召回Flink 实时结果做加权推荐网关按时间衰减合并两路结果。2.2 实时推荐的数据输入行为流长什么样埋点行为是最常见的数据源一条完整记录大致长这样{ userId: 12345, skuId: 8765, behaviorType: click, scene: home_feed, ts: 2025-01-12T14:23:11Z }userId是用户标识skuId是商品标识behaviorType区分点击、收藏、加购、下单和曝光。scene标记推荐位比如首页信息流、搜索结果页、详情页相关推荐。ts是行为发生时间推荐系统必须用它作为事件时间不能用 Flink 收到消息的处理时间替代否则消息在 Kafka 里积压时会直接把时间语义搞乱。曝光数据容易被人忽略它其实很重要推荐结果中用户没点的商品就是天然负样本。做实时排序模型时曝光流和点击流需要按用户、商品、场景做关联这一步也放在 Flink 里完成。2.3 Flink的三个机制决定实时推荐可行性Flink 能承担实时推荐靠的是事件时间、状态和精确一次处理三个机制。代码里落地事件时间是这样的DataStreamUserBehavior stream env .addSource(kafkaSource) .assignTimestampsAndWatermarks( WatermarkStrategy.UserBehaviorforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, ts) - event.ts.toEpochMilli()) );forBoundedOutOfOrderness(Duration.ofSeconds(10))表示允许行为事件最多乱序 10 秒超过这个延迟的数据算迟到数据。这个值不是越小越好设成 0 会大量丢弃乱序数据设成 60 秒又会拖延窗口触发。推荐场景一般 10 到 20 秒比较合理既保证实时性又能容忍用户手机网络抖动造成的乱序。状态机制也关键。商品相似度计算要在内存里维护“每个用户最近看了哪些商品”用户的实时偏好向量也要持续更新。Flink 的ValueState、MapState会把中间数据存在本地配合 checkpoint 定期持久化算子挂掉后能从最近一次 checkpoint 恢复不会把中间计数全部归零。2.4 为什么不是Storm或Spark StreamingStorm 延迟确实低但状态管理和精确一次语义支持弱实现一个“用户最近浏览序列”要自己管理外部存储工程成本高。Spark Streaming 的微批模型让事件时间窗口别别扭扭调 watermark 和迟到数据的处理都绕。Flink 的 DataStream API 和 SQL 双 API 设计则实用得多简单特征计算用 SQL 表达复杂的状态逻辑用 DataStream 实现一个平台都能覆盖。这也是实时推荐选 Flink 而不是其他框架的理由。3. 商品实时推荐系统的Flink链路与组件选型实时推荐系统不是只有一个 Flink 作业而是多条计算链路的组合。设计架构时我习惯把链路按方向拆开行为进、特征算、候选出、服务查。3.1 从埋点到推荐结果的五段链路链路整体走向是客户端埋点 → Kafka → Flink 计算 → Redis/ES → 推荐服务聚合。埋点数据先进 Kafka原因很简单削峰填谷。用户行为有明显的波峰波谷早晚高峰和促销活动时期的流量能差好几倍。Kafka 把流量缓冲住Flink 按自己的节奏消费推荐服务不会被打垮。Flink 在中间承担三类计算实时特征、实时相似度、实时候选生成。实时特征包括点击次数、加购次数、浏览时长、最近浏览序列等短期信号实时相似度指商品与商品的关联强度用共现或协同过滤思路算实时候选生成是把用户短期行为和相似商品结合产出候选列表。算完的结果落到 Redis 和 ES。Redis 存用户实时向量和候选列表KV 查询快ES 存最终推荐结果方便按条件过滤和做 AB 实验。推荐服务读取这些结果后再和离线模型结果合并排序输出。3.2 存储层选型Redis、ES还是ClickHouse存储角色推荐系统里承担的任务选型注意RedisKV 存储用户实时偏好、TopN 候选、特征缓存设置合理 TTL防 key 无限增长ES文档存储最终推荐结果、商品索引、结果查询批量写入避免逐条写入打满节点ClickHouse列式分析行为明细归档、离线验证、特征回填适合 OLAP 查询不适合高频点查我一般把 Redis 当第一道查询缓存推荐服务先用 userId 查候选列表。没命中再去 ES 查。如果用户量巨大可以再在前面加一层本地缓存但要容忍几秒钟的数据延迟。3.3 三种推荐任务在Flink里的拆分方式实时特征、实时召回、实时排序不建议写进同一个作业里。一个作业承担太多逻辑任何人都得小心维护。常见做法是拆三个作业作业技术选型输出典型窗口实时特征作业Flink SQL用户特征明细5 分钟滚动窗口实时召回作业DataStream API商品相似度、候选列表7 天滑动窗口实时过滤作业Flink SQL CDC 维表过滤下架、已购商品无窗口逐条过滤实时特征作业消耗资源小可以用 Flink SQL 快速迭代实时召回作业状态大用 DataStream API 精细控制状态和触发时机实时过滤作业要感知商品上下架状态配合 CDC 维表实现。拆开后每个作业独立重启、独立扩容模型迭代不会互相拖累。3.4 工程化的Flink代码目录与启动骨架刚接触 Flink 的开发者常把所有逻辑写进一个 main 方法。这样也能跑但维护三个作业时会非常痛苦。一个相对工程化的目录结构长这样realtime-rec ├── src/main/java/com/example/rec │ ├── Main.java │ ├── connector │ │ ├── KafkaSinkProvider.java │ │ ├── RedisSinkProvider.java │ │ └── EsSinkProvider.java │ ├── function │ │ ├── RecentItemsFunction.java │ │ ├── CoOccurrenceEmitter.java │ │ └── SimilarityAccumulator.java │ ├── job │ │ ├── FeatureJob.java │ │ └── RecallJob.java │ └── config │ └── ParamInitializer.java ├── sql │ ├── feature_job.sql │ └── recall_filter.sql └── deploy └── submit_job.shjob包里的每个类对应一个可提交的作业入口connector包统一管理 sink 创建逻辑sql目录存可执行的 Flink SQL 脚本。这样做的直接好处是新同学接手时看目录就能明白作业边界在哪。Main 方法的骨架一般长这样public static void main(String[] args) { ParameterTool params ParameterTool.fromArgs(args); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(params.getInt(checkpoint.interval, 60_000)); env.setStateBackend(new HashMapStateBackend()); String kafkaBrokers params.getRequired(kafka.brokers); String redisHost params.getRequired(redis.host); // 所有连接参数从命令行传入不硬编码 ... }ParameterTool.fromArgs可以从启动参数读取配置checkpoint.interval默认 60 秒生产环境一般设置在 60 到 120 秒之间。连接信息必须参数化不同环境测试、预发、生产直接换启动脚本而不是改代码重新打包。硬编码连接串是生产事故的高发源头。4. 用Flink消费Kafka、算相似度并写入目标端链路设计和选型只能算纸上谈兵这章把核心代码跑起来从 Kafka 接入行为流实时算商品相似度再算用户特征最后写入 Redis 和 ES。4.1 Flink消费Kafka的最小可运行配置Flink SQL 里接入 Kafka 行为流一张 DDL 就能搞定CREATE TABLE kafka_behaviors ( userId BIGINT, skuId BIGINT, behaviorType STRING, scene STRING, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector kafka, topic user-behavior, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, properties.group.id realtime-rec, scan.startup.mode earliest-offset, format json, json.timestamp-format.standard ISO-8601 );WATERMARK FOR ts AS ts - INTERVAL 5 SECOND表示允许 5 秒乱序延迟和上一章的 10 秒策略相互呼应。scan.startup.mode支持earliest-offset和latest-offset离线补数测试用前者线上接入用后者更常见。properties.group.id决定 Kafka 消费位点的管理粒度变更 group id 会从头或从最新位点重新消费生产环境要谨慎修改。4.2 实时商品相似度计算的DataStream实现实时 ItemCF 的基本逻辑是同一用户的行为序列里两件商品出现得越近相似度越高。核心算子链如下DataStreamItemPair pairs behaviors .keyBy(b - b.userId) .process(new RecentItemsFunction()) .flatMap(new CoOccurrenceEmitter()); pairs.keyBy(pair - pair.itemA _ pair.itemB) .process(new SimilarityAccumulator()) .addSink(redisSink);RecentItemsFunction维护每个用户最近购买的 30 个商品CoOccurrenceEmitter把当前商品和序列里的历史商品两两组成一对。SimilarityAccumulator对共现次数增量累加周期输出到 Redis。其中SimilarityAccumulator的实现是public static class SimilarityAccumulator extends KeyedProcessFunctionString, ItemPair, SimilarityResult { private ValueStateLong countState; Override public void open(Configuration parameters) { countState getRuntimeContext().getState(new ValueStateDescriptor(pair-count, Long.class)); } Override public void processElement(ItemPair pair, Context ctx, CollectorSimilarityResult out) throws Exception { Long count countState.value(); count count null ? 1L : count 1L; countState.update(count); ctx.timerService().registerProcessingTimeTimer( ctx.timerService().currentProcessingTime() 60_000); } Override public void onTimer(long timestamp, OnTimerContext ctx, CollectorSimilarityResult out) throws Exception { out.collect(new SimilarityResult(currentKey(), countState.value())); countState.clear(); } }这段代码用ValueState保存每对商品的共现次数数据在算子本地维护不产生外部 IO。定时器每隔 60 秒触发一次输出当前累计值后清零。生产环境建议把处理时间定时器换成事件时间窗口加自定义 Trigger能更好对齐业务时间语义。4.3 用Flink SQL计算用户短期行为特征用户特征用 SQL 表达更清晰。下面这段计算每个用户最近 5 分钟的点击量和点击商品数INSERT INTO user_clicks_5m SELECT userId, COUNT(*) AS click_cnt, COUNT(DISTINCT skuId) AS distinct_skus, TUMBLE_START(ts, INTERVAL 5 MINUTE) AS win_start FROM kafka_behaviors WHERE behaviorType click GROUP BY userId, TUMBLE(ts, INTERVAL 5 MINUTE);TUMBLE(ts, INTERVAL 5 MINUTE)是滚动窗口每 5 分钟输出一次结果。TUMBLE_START取出窗口起始时间用于下游关联时标记特征更新时间。窗口时长是这里的关键参数点击类特征用 5 分钟加购和下单类信号稀疏建议用 1 小时窗口否则大部分窗口是空值。真实推荐链路里还会关联商品维表补全类目和价格信息SELECT b.userId, b.skuId, d.categoryId FROM kafka_behaviors b JOIN dim_sku FOR SYSTEM_TIME AS OF b.ts AS d ON b.skuId d.skuId;FOR SYSTEM_TIME AS OF b.ts是维表关联的固定写法系统会按事件时间取出当时的维表快照。关联操作用来过滤已经下架的商品也可以给行为数据补上商品类目用于算类目级别的用户偏好。4.4 推荐结果写入ES的SQL实现结果写入 ES用 SQL 指定 connector 即可CREATE TABLE sink_recall ( userId BIGINT, skuList STRING, updateTime TIMESTAMP(3), PRIMARY KEY (userId) NOT ENFORCED ) WITH ( connector elasticsearch-7, index recall_result, sink.bulk-flush.max-actions 500, sink.bulk-flush.max-size 10mb, sink.bulk-flush.interval 5000, format json );sink.bulk-flush.max-actions控制攒够 500 条批量写入一次sink.bulk-flush.max-size是 10MB 触发一次sink.bulk-flush.interval兜底每 5 秒强制刷一次。这三个参数必须调直接用默认值会在低流量时一条一条写高流量时把 ES 打满。PRIMARY KEY (userId) NOT ENFORCED不是给数据库建主键而是让 ES sink 用userId做文档 id实现 upsert 语义。Redis 一般用 DataStream 自定义 sink 写入结构是rec:user:{userId}存最近购买商品 ID 列表TTL 设置 24 小时即可。注意 Redis 缓存和 ES 的最终一致性问题先写 ES 再更新 Redis推荐服务短暂读到旧数据可以接受但反过来 Redis 先更新、ES 失败会导致明细不一致排查时更费时间。4.5 推荐线路的幂等与一致性设计实时链路里 sink 重试是常态幂等不是可选项。ES 主键天然幂等Redis 的 set 操作幂等MySQL 写入必须指定主键并做 upsert。Kafka 开启 checkpoint 之后作业重启会重复消费部分消息下游必须能接受重复写入。这里的原则是不要依赖“只写一次”要保证“重复写结果一致”。5. Flink实时推荐的参数调优与连接器异常排查参数调优和异常排查在实时推荐系统里占的工时比重很大。数据源、窗口、连接器每一层都有各自的坑。5.1 watermark不推进时SQL任务的表现和排查现象是结果表很久不出数据作业没有反压日志里也看不到异常。这种问题十有八九出在事件时间上。在 SQL Client 里跑一条窗口查询看窗口能不能按时触发SELECT TUMBLE_ROWTIME(ts, INTERVAL 5 MINUTE) AS win_end, COUNT(*) AS total_cnt FROM kafka_behaviors WHERE behaviorType click GROUP BY TUMBLE(ts, INTERVAL 5 MINUTE);如果这个查询 5 分钟不产出结果说明 watermark 没有推进。Kafka 多分区场景下watermark 取所有分区的最小值。某个分区没有新消息整个 watermark 就卡住。DataStream API 里可以用空闲分区机制解决WatermarkStrategy .UserBehaviorforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withIdleness(Duration.ofMinutes(2))withIdleness(Duration.ofMinutes(2))表示分区 2 分钟没有新数据时把该分区标记为空闲不再参与最小水位计算。但 Flink SQL 的 Kafka connector 没有直接暴露这个参数生产上我一般靠监控 topic 各分区 LAG 和在上游写入心跳消息来解决后者实践效果更稳定。5.2 Flink JDBC连接器异常的常见症状“Flink的JDBC连接器异常”这个关键词铺开来就是一张问题排查表报错信息常见根因处置方向Connection is not available, request timed out连接池被打满开启 lookup.cache增加最大连接数Communications link failureMySQL 主动断开空闲连接调大 wait_timeout启用连接校验Connection reset by peer大量异常重启导致连接残留检查作业重启频率设置重试上限最常见的是维表 join 不设缓存每条行为数据都到 MySQL 查一次维表高峰期瞬间把连接池打满。JDBC 维表 DDL 里必须配置缓存CREATE TABLE dim_sku ( skuId BIGINT, skuName STRING, categoryId BIGINT, PRIMARY KEY (skuId) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://mysql-host:3306/rec, table-name sku_dim, lookup.cache.max-rows 10000, lookup.cache.ttl 30min, lookup.max-retries 3 );lookup.cache.max-rows设 10000 行lookup.cache.ttl设 30 分钟。缓存是为了削峰热门商品的维表数据在 TM 内存里直接命中只有冷门商品才穿透到 MySQL。代价是数据更新最多延迟 30 分钟商品改名、类目调整这类低频变更可以接受。5.3 Flink CDC同步商品状态作为实时过滤维表推荐结果不能推荐已下架商品。商品上下架状态存在业务 MySQL 里用 Flink CDC 同步成实时维表是标准方案CREATE TABLE dim_sku_cdc ( skuId BIGINT PRIMARY KEY NOT ENFORCED, skuName STRING, categoryId BIGINT, status INT, updateTime TIMESTAMP(3) ) WITH ( connector mysql-cdc, hostname mysql-host, port 3306, username rec, password ***, scan.startup.mode initial, server-id 5400-5404 );scan.startup.mode initial表示启动时先做一次全量快照再增量消费 binlog。server-id的范围设置很重要每个并行度需要独立的 server-id如果并行度是 5这里至少配 5 个连续 ID。配少了会出现 binlog 连接冲突。CDC 任务的核心价值是让推荐过滤逻辑跟着业务变化走商品下架几秒后推荐链路就能把它从候选里剔除而不是等离线任务第二天刷新。5.4 状态后端、checkpoint与standalone部署调优实时相似度作业的状态会持续增长状态后端选择直接影响稳定性。状态后端适用场景生产注意HashMapStateBackend状态小于 1GB内存不足大状态 OOM 风险高RocksDBStateBackend大状态超过可用内存磁盘 IO 有开销启用增量 checkpointRocksDB 增量快照超大状态生产环境每个 TM 管理内存上限execution.checkpointing.interval60000是生产常用间隔。RocksDB 模式下建议开启state.backend.rocksdb.memory.managedtrue让 Flink 自动管理 RocksDB 使用的堆外内存避免和 JVM 堆内存互相争抢。standalone 部署模式下配置taskmanager.memory.process.size时要为系统进程预留 2GB 左右别把物理内存全分给 Flink否则 Flink 一压测节点直接 OOM。5.5 反压定位的三步操作实时推荐任务出问题一个典型场景是用户流量突然上涨。定位步骤基本固定打开 Flink Web UI 的 BackPressure 页签找到 High 状态的算子。反压在 source 和下游算子之间多半是读取速度超过处理速度检查下游 sink 的批量参数。反压在 keyBy 和 window 附近大概率是状态太大或数据倾斜查看大 key 分布必要时对 key 加盐二次聚合。线上任务最好把反压监控接入告警平台连续 5 分钟处于 HIGH 就通知值班人员。推荐链路实时性要求高等到用户反馈“推荐不更新”再排查损失已经产生。6. 用Flink SQL Gateway做实时推荐的旁路验证最后一个环节说验证。实时推荐任务上线后不要直接信任 sink 里的数据要能随时从旁路确认链路是通的。6.1 在SQL Client或SQL Gateway里做旁路查询线上作业出问题优先在 SQL Gateway 里起一个临时查询复用线上的 connector 配置查同一张 Kafka 表。这样做的好处是不影响线上作业又能直接看到当前流里的数据长什么样SELECT userId, skuId, behaviorType, CURRENT_TIMESTAMP AS query_time FROM kafka_behaviors WHERE ts CURRENT_TIMESTAMP - INTERVAL 5 MINUTE LIMIT 10;如果能查到数据说明 Kafka 接入和 schema 解析正常。如果查不到优先检查scan.startup.mode是不是被改过、group id 是否换过以及 watermark 是否卡住。用这条 SQL 能把问题范围快速缩小到“链路”还是“计算逻辑”。6.2 注入测试行为验证端到端链路延迟旁路查询只能证明数据进来了端到端延迟得靠注入一条测试行为来验证。用 Kafka 自带的生产者工具发送一条测试消息echo {userId: 999001, skuId: 777701, behaviorType: click, scene: home_feed, ts: 2025-01-12T14:23:11Z} | \ kafka-console-producer.sh --broker-list kafka-1:9092 --topic user-behavior然后立刻查看 Redis 里对应的用户特征缓存是否更新redis-cli GET rec:user:999001窗口类特征有固定计算周期5 分钟窗口的任务最多等 5 分钟零几秒如果超过 8 分钟还没更新基本可以断定窗口触发或 sink 写入选型有问题。这个方法对新手特别友好不打断线上任务不依赖复杂监控手动操作一遍就能确认主链路健康。6.3 为特征表补一份轻量血缘记录数据血缘在实时链路里常被忽略等排错时才想起来要查“这个字段哪来的”。Flink SQL 的EXPLAIN能看到算子的输入输出关系但不会自动沉淀成跨任务的血缘文档。我一般会在 sql 目录里给每张结果表维护一段注释-- table: user_clicks_5m -- source: kafka_behaviors -- filter: behaviorType click -- window: tumble 5 minute -- sink: es recall_result index注释写清楚来源表、过滤条件、窗口类型和下游消费方排错时先看注释再翻代码效率高一截。真正的全局血缘需要前端采集 catalog 信息和 SQL 解析结果那是平台团队的事但在小团队里一份维护良好的 SQL 注释已经能解决八成追溯需求。Flink SQL Gateway 配合 DDL 注释基本能满足实时推荐链路日常运维的数据发现需求。本文还有配套的精品资源点击获取