资讯动态

Flink+Kafka实时推荐系统:评分行为驱动与离线兜底实战

发布时间:2026/10/6 2:55:36 来源:尧图企业网站定制
简介本资源为基于Flink的商品实时推荐系统完整项目源码面向计算机、人工智能、大数据等相关专业的在校学生、教师及企业开发人员可用于毕业设计、课程设计、项目立项演示或技术进阶学习。项目以Kafka作为数据通道当用户产生评分行为时将数据发送至Flink结合用户历史评分实现实时与离线双路推荐实时推荐涵盖基于行为和实时热门离线推荐则包括历史热门、历史优质商品及ItemCF协同过滤完整覆盖推荐系统核心链路。资源包共408个文件以xml配置、java源码、class编译文件为主辅以ts、vue、stylus等前端资源及properties、csv、sql、json等数据与配置文档压缩包约4.27MB目录结构清晰便于按模块检索与二次开发。目前已有87人学习下载。项目代码经过测试运行成功并获导师指导认可答辩评审95分读者可据此掌握Flink实时计算与推荐算法整合思路也可在现有代码基础上修改扩展功能。1. 评分行为驱动的实时推荐Flink 与 Kafka 到底在链路里扛了什么用户在前端点下一颗星这条评分行为从产生到影响推荐结果中间要穿过消息队列、流处理引擎、特征存储和召回排序四层。很多团队一开始把这件事想简单了不就是把评分写进数据库然后定时跑一遍协同过滤吗真到线上就会发现离线批处理跑一次要几十分钟用户刚评完一部电影推荐列表纹丝不动体验割裂得厉害。基于 Flink 的商品实时推荐系统核心就是让评分行为在秒级内进入计算链路实时更新用户兴趣向量同时保留离线全量训练做兜底。Kafka 在这里承担削峰和解耦Flink 负责有状态流计算两者配合才能把“实时推荐”从概念落到可运维的工程。这套方案适合已经有离线推荐基础、想补上实时通路的团队也适合刚接触流式计算、想找一个完整场景练手的工程师。下面按数据怎么流、状态怎么管、离线怎么补、坑怎么避的顺序拆开讲。2. 从评分事件到 Kafka Topic数据接入的格式约定与分区策略2.1 评分事件该长什么样实时推荐链路的第一道关卡是数据格式。评分行为看起来简单——用户 ID、商品 ID、评分值、时间戳但真到生产环境字段缺失、时间戳格式混乱、评分值越界是家常便饭。我一般会在 Kafka 消息体里用 JSON 承载字段固定为user_id、item_id、rating、event_time、trace_id。trace_id不是必须的但排查问题时没有它你只能靠猜。event_time用毫秒级 Unix 时间戳不要用字符串否则 Flink 里做事件时间窗口时还得额外解析容易翻车。{ user_id: u10023, item_id: i8871, rating: 4.5, event_time: 1713427200000, trace_id: req-9f3a2b }评分值范围建议在接入层就做校验比如限制在 0.5 到 5.0 之间超出范围直接丢到死信 Topic不要让它污染下游状态。死信 Topic 的命名可以用rating_dlq保留原始消息和失败原因方便后续补录。2.2 Kafka Topic 的分区与保留策略Topic 分区数直接决定 Flink 算子的并行度上限。如果分区数只有 3Flink 的 Source 并行度开到 6 也没用多出来的 subtask 会空转。我一般按峰值 QPS 来估算单分区吞吐按 5MB/s 到 10MB/s 算评分事件单条大约 200 字节峰值 1 万 QPS 的话3 到 6 个分区足够。但要注意分区数一旦确定后续扩分区会导致同一用户的数据散落到不同分区如果 Flink 里按用户 ID 做 keyBy扩分区后状态迁移会很麻烦。所以初期宁可多分几个比如 12 个分区后面用不上也比不够强。保留策略方面评分事件建议保留 7 天。实时链路消费完就可以提交 offset但离线训练可能需要回放最近几天的数据保留太短会导致离线补数据时无源可查。Kafka 的retention.ms设成604800000同时cleanup.policy用delete而不是compact因为评分事件是追加型事实数据不需要保留每个 key 的最新值。# 创建评分事件 Topic12 分区副本 2保留 7 天 kafka-topics.sh --create \ --bootstrap-server localhost:9092 \ --topic rating_events \ --partitions 12 \ --replication-factor 2 \ --config retention.ms604800000 \ --config cleanup.policydelete副本数至少 2单副本在 broker 重启时会丢数据。如果集群规模允许3 副本更稳但写入延迟会略高。生产环境建议min.insync.replicas2配合acksall确保消息写入至少两个副本才返回成功。2.3 生产者端的可靠性配置评分事件的生产者通常是后端服务用 Kafka 客户端直接发。关键配置有三个acks、retries、enable.idempotence。acksall保证消息被所有同步副本确认retries设成Integer.MAX_VALUE配合delivery.timeout.ms控制重试总时长enable.idempotencetrue防止重试导致消息重复。这三个一起开基本能保证不丢不重。Properties props new Properties(); props.put(bootstrap.servers, kafka1:9092,kafka2:9092,kafka3:9092); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(acks, all); props.put(retries, Integer.MAX_VALUE); props.put(enable.idempotence, true); props.put(max.in.flight.requests.per.connection, 5); props.put(delivery.timeout.ms, 120000);max.in.flight.requests.per.connection在开启幂等性后Kafka 2.0 以上版本可以设到 5既保证顺序又提高吞吐。如果设成 1吞吐会明显下降但顺序性最强。评分事件对顺序不敏感用 5 就行。3. Flink 实时计算用 KeyedProcessFunction 维护用户兴趣向量3.1 为什么选 KeyedProcessFunction 而不是窗口聚合实时推荐的核心逻辑是每来一条评分就更新该用户的兴趣向量然后基于新向量生成推荐。窗口聚合适合做统计类指标比如“最近 10 分钟平均评分”但推荐需要的是逐条更新、立即生效。KeyedProcessFunction提供了processElement和onTimer两个入口前者处理每条评分后者可以做超时清理或定时衰减。按user_id做 keyBy 后同一用户的所有评分都会落到同一个 subtask状态天然隔离。DataStreamRatingEvent ratings env .addSource(new FlinkKafkaConsumer(rating_events, new JSONDeserializationSchema(), props)) .assignTimestampsAndWatermarks( WatermarkStrategy.RatingEventforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, ts) - event.getEventTime()) ); DataStreamRecommendation recommendations ratings .keyBy(RatingEvent::getUserId) .process(new UserInterestProcessFunction());forBoundedOutOfOrderness(Duration.ofSeconds(5))表示允许 5 秒乱序超过 5 秒的事件会被丢弃或进入侧输出流。评分行为通常不会乱序太久5 秒足够。如果 Kafka 分区之间时钟差异大可以放宽到 10 秒但会增加状态驻留时间。3.2 用户兴趣向量的状态结构用户兴趣向量我一般用MapStateString, Double存key 是商品 ID 或品类 IDvalue 是兴趣权重。每来一条评分按公式更新新权重 旧权重 × 衰减因子 评分值 × 即时权重。衰减因子取 0.95 到 0.99 之间即时权重取 0.1 到 0.3。这样既能保留历史兴趣又能让近期行为快速影响推荐。public class UserInterestProcessFunction extends KeyedProcessFunctionString, RatingEvent, Recommendation { private MapStateString, Double interestState; private static final double DECAY 0.97; private static final double IMMEDIATE_WEIGHT 0.2; Override public void open(Configuration parameters) { MapStateDescriptorString, Double descriptor new MapStateDescriptor(interest, String.class, Double.class); interestState getRuntimeContext().getMapState(descriptor); } Override public void processElement(RatingEvent event, Context ctx, CollectorRecommendation out) throws Exception { // 先对所有已有兴趣做衰减 for (Map.EntryString, Double entry : interestState.entries()) { interestState.put(entry.getKey(), entry.getValue() * DECAY); } // 再更新当前商品的兴趣权重 String itemKey event.getItemId(); double current interestState.contains(itemKey) ? interestState.get(itemKey) : 0.0; interestState.put(itemKey, current event.getRating() * IMMEDIATE_WEIGHT); // 注册定时器30 分钟后清理低权重兴趣 ctx.timerService().registerProcessingTimeTimer(ctx.timerService().currentProcessingTime() 30 * 60 * 1000); // 基于更新后的向量生成推荐 ListString topItems interestState.entries().stream() .sorted((a, b) - Double.compare(b.getValue(), a.getValue())) .limit(10) .map(Map.Entry::getKey) .collect(Collectors.toList()); out.collect(new Recommendation(event.getUserId(), topItems, System.currentTimeMillis())); } Override public void onTimer(long timestamp, OnTimerContext ctx, CollectorRecommendation out) throws Exception { // 清理权重低于阈值的兴趣防止状态无限膨胀 IteratorMap.EntryString, Double it interestState.entries().iterator(); while (it.hasNext()) { if (it.next().getValue() 0.01) { it.remove(); } } } }DECAY和IMMEDIATE_WEIGHT这两个参数需要根据业务调。评分行为稀疏的场景衰减可以慢一点比如 0.99评分频繁的场景衰减快一点0.95 更能反映近期兴趣。onTimer里的清理阈值 0.01 是经验值低于这个值的兴趣对推荐结果几乎没有影响留着只会拖慢状态遍历。3.3 状态后端与 Checkpoint 配置Flink 的状态默认存在 TaskManager 堆内存里用户量一大就会 OOM。生产环境建议用 RocksDB 状态后端把状态落到本地磁盘堆内存只保留热数据。Checkpoint 间隔设成 1 到 3 分钟太短会影响吞吐太长会导致故障恢复时重放大量数据。env.setStateBackend(new RocksDBStateBackend(hdfs://namenode:8020/flink/checkpoints)); env.enableCheckpointing(120000); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(60000); env.getCheckpointConfig().setCheckpointTimeout(300000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);setMinPauseBetweenCheckpoints(60000)确保两次 Checkpoint 之间至少间隔 1 分钟避免 Checkpoint 频繁触发拖垮正常处理。setMaxConcurrentCheckpoints(1)防止多个 Checkpoint 并发执行RocksDB 状态下并发 Checkpoint 容易导致文件句柄耗尽。4. 离线推荐兜底用 Flink SQL 把评分行为同步到 Hive 与 Redis4.1 离线链路的定位与数据流向实时推荐解决的是“秒级响应”但冷启动用户和新商品没有足够行为数据实时向量几乎是空的。离线推荐用全量历史评分跑协同过滤或矩阵分解每天更新一次覆盖实时链路的盲区。离线链路的数据源也是 Kafka 里的评分事件但消费方式不同实时链路用 Flink DataStream 逐条处理离线链路用 Flink SQL 做微批写入 Hive再跑 Spark 或 Hive SQL 训练模型。-- Flink SQL 将 Kafka 评分事件写入 Hive 离线表 CREATE TABLE rating_kafka ( user_id STRING, item_id STRING, rating DOUBLE, event_time BIGINT, ts AS TO_TIMESTAMP_LTZ(event_time, 3), WATERMARK FOR ts AS ts - INTERVAL 10 SECOND ) WITH ( connector kafka, topic rating_events, properties.bootstrap.servers kafka1:9092,kafka2:9092, properties.group.id offline-rating-group, scan.startup.mode earliest-offset, format json ); CREATE TABLE rating_hive ( user_id STRING, item_id STRING, rating DOUBLE, event_time BIGINT, dt STRING ) PARTITIONED BY (dt) WITH ( connector hive, hive-version 3.1.2, default-database recommend, sink.partition-commit.policy.kind metastore,success-file ); INSERT INTO rating_hive SELECT user_id, item_id, rating, event_time, DATE_FORMAT(TO_TIMESTAMP_LTZ(event_time, 3), yyyy-MM-dd) AS dt FROM rating_kafka;scan.startup.mode用earliest-offset表示从最早的消息开始消费适合离线补全量数据。如果只需要增量改成latest-offset或group-offsets。sink.partition-commit.policy.kind设成metastore,success-file确保 Hive 元数据更新和成功标记文件同时写入否则下游 Spark 读不到新分区。4.2 离线推荐结果写回 Redis 供实时链路查询离线训练出的推荐结果通常是一批用户对一批商品的预测评分数据量不大适合写进 Redis 做在线查询。实时链路生成推荐时如果用户兴趣向量为空或商品候选不足就从 Redis 拉离线结果兜底。写入 Redis 可以用 Flink SQL 的 Redis Connector也可以用离线任务跑完后用脚本批量导入。CREATE TABLE offline_rec_redis ( user_id STRING, item_id STRING, score DOUBLE, PRIMARY KEY (user_id, item_id) NOT ENFORCED ) WITH ( connector redis, host redis-host, port 6379, redis-mode single, command SET, key-column user_id, value-column item_id,score, expire 86400 ); INSERT INTO offline_rec_redis SELECT user_id, item_id, score FROM offline_recommendation_result;expire设成 86400 秒即一天过期和离线任务的更新周期对齐。如果离线任务延迟Redis 里的旧数据还能撑一段时间不至于完全不可用。key-column用user_idvalue 里拼item_id和score查询时用GET user_id拿到一个字符串再解析。如果推荐列表较长可以用 Redis 的 List 或 ZSet但 Flink SQL Connector 对复杂结构的支持有限简单字符串拼接更稳。4.3 实时与离线的融合策略实时推荐和离线推荐不是二选一而是按场景融合。用户有近期行为时实时向量权重高用户行为稀疏时离线结果权重高。我一般用一个简单的加权公式最终得分 α × 实时得分 (1 - α) × 离线得分α 根据用户最近 7 天的行为条数动态调整。行为条数超过 20 条α 取 0.8少于 5 条α 取 0.2。这个逻辑可以放在 Flink 的processElement里也可以放在推荐服务层做。double alpha Math.min(0.8, Math.max(0.2, recentBehaviorCount / 25.0)); double finalScore alpha * realtimeScore (1 - alpha) * offlineScore;recentBehaviorCount可以从 Flink 状态里维护一个计数器每来一条评分加一每天零点重置。这个计数器用ValueStateLong存配合定时器做日切。5. 避坑与排查评分行为实时推荐链路的 5 个血泪教训5.1 现象Flink 任务频繁重启Checkpoint 一直失败原因通常是 RocksDB 状态太大Checkpoint 超时。用户兴趣向量如果每个用户存几千个商品百万用户就是几十 GB 状态RocksDB 写 HDFS 一次要几分钟。解决方法是给状态设 TTLFlink 的StateTtlConfig可以自动清理过期状态。StateTtlConfig ttlConfig StateTtlConfig.newBuilder(Time.days(7)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .cleanupInRocksdbCompactFilter(1000) .build(); MapStateDescriptorString, Double descriptor new MapStateDescriptor(interest, String.class, Double.class); descriptor.enableTimeToLive(ttlConfig);cleanupInRocksdbCompactFilter(1000)表示每处理 1000 条记录触发一次清理值越小清理越频繁但开销越大。7 天 TTL 对推荐场景够用超过 7 天没行为的用户兴趣本来就应该衰减掉。5.2 现象Kafka 消息延迟高Flink 消费跟不上先看 Kafka 的消费者 lag如果 lag 持续增长说明 Flink 处理能力不足。常见原因是keyBy后数据倾斜某个用户的行为特别多落到一个 subtask 上。解决方法是给 key 加盐比如user_id _ (hash % 4)把热点用户拆到多个 subtask但这样同一用户的状态会分散需要在推荐时做二次聚合。另一种方法是提高 Flink 并行度但前提是 Kafka 分区数够。// 加盐打散热点 key DataStreamRatingEvent salted ratings .map(event - { int salt Math.abs(event.getUserId().hashCode()) % 4; event.setSaltedKey(event.getUserId() _ salt); return event; }) .keyBy(RatingEvent::getSaltedKey);加盐后同一用户的评分会分散到 4 个 subtask每个 subtask 维护部分兴趣向量。生成推荐时需要把 4 个部分向量合并可以在 Flink 里再做一次keyBy(user_id)聚合或者写到外部存储后由推荐服务合并。5.3 现象离线 Hive 表分区没有数据Spark 读不到Flink SQL 写 Hive 时sink.partition-commit.policy.kind如果只配了success-file没配metastoreHive 元数据不会更新Spark 查表时看不到新分区。必须两个都配。另外sink.partition-commit.delay默认是 0如果数据还没写完就提交分区会导致分区里只有部分文件。建议设成1min等数据落稳再提交。sink.partition-commit.delay 1min, sink.partition-commit.policy.kind metastore,success-file如果 Hive 表是 ORC 格式还要确认 Flink 的 Hive Connector 版本和 Hive 版本匹配版本不匹配会报ClassNotFoundException这个坑我踩过两次。5.4 现象Redis 里离线推荐结果被实时结果覆盖实时链路和离线链路都往 Redis 写同一个 key实时链路写的是临时推荐离线链路写的是全量推荐互相覆盖导致推荐列表忽长忽短。解决方法是分开 key 前缀实时用realtime:rec:{user_id}离线用offline:rec:{user_id}推荐服务查询时合并两个来源。或者用 Redis 的 Hash 结构实时和离线写不同 field读取时HGETALL一起拿。5.5 现象评分事件时间戳乱序导致推荐结果抖动用户可能先评了 A 再评 B但 B 的消息先到 Flink。如果没设水位线Flink 按处理时间处理B 先更新兴趣向量A 后到又覆盖回去推荐结果就会抖。必须用事件时间加水印forBoundedOutOfOrderness设 5 到 10 秒让 Flink 等一等迟到数据。如果乱序超过水位线数据会被丢弃可以配侧输出流把迟到数据收集起来后续补处理。OutputTagRatingEvent lateTag new OutputTagRatingEvent(late-ratings){}; SingleOutputStreamOperatorRecommendation mainStream ratings .keyBy(RatingEvent::getUserId) .process(new UserInterestProcessFunction()) .getSideOutput(lateTag);侧输出流的数据可以写到另一个 Kafka Topic离线任务再消费一次保证不丢。6. 进阶技巧用 Flink 的 Async I/O 查 Redis 补全商品特征实时推荐生成候选商品后通常还需要补全商品特征比如品类、价格、销量这些数据存在 Redis 或 HBase 里。如果直接在processElement里同步查 Redis每条评分都要等一次网络往返吞吐上不去。Flink 的 Async I/O 可以并发查外部存储吞吐能提升几倍。DataStreamRecommendation enriched AsyncDataStream .unorderedWait( recommendations, new AsyncRedisLookup(), 5000, // 超时 5 秒 TimeUnit.MILLISECONDS, 100 // 最大并发请求数 ); public class AsyncRedisLookup extends RichAsyncFunctionRecommendation, Recommendation { private transient JedisPool jedisPool; Override public void open(Configuration parameters) { jedisPool new JedisPool(redis-host, 6379); } Override public void asyncInvoke(Recommendation rec, ResultFutureRecommendation resultFuture) { CompletableFuture.supplyAsync(() - { try (Jedis jedis jedisPool.getResource()) { for (String itemId : rec.getItems()) { String feature jedis.get(item:feature: itemId); rec.addFeature(itemId, feature); } return rec; } }).thenAccept(resultFuture::complete); } }unorderedWait表示不保证结果顺序吞吐更高如果推荐列表对顺序敏感用orderedWait但会牺牲一些性能。5000毫秒超时和100并发数需要根据 Redis 的承载能力调并发太高会把 Redis 打满。我一般先设 50压测后逐步往上加同时监控 Redis 的latency和connected_clients。另一个技巧是给 Redis 查询加本地缓存。商品特征变化不频繁可以在RichAsyncFunction的open方法里初始化一个 Guava Cache设置 5 分钟过期减少对 Redis 的重复查询。这样即使 Redis 抖动推荐链路也能靠本地缓存撑一段时间。private transient CacheString, String localCache; Override public void open(Configuration parameters) { localCache CacheBuilder.newBuilder() .maximumSize(10000) .expireAfterWrite(5, TimeUnit.MINUTES) .build(); }本地缓存和 Async I/O 配合Redis 的 QPS 能降一个数量级。但要注意缓存过期时间不能太长否则商品下架后推荐里还会出现。5 分钟是个折中值对大多数电商场景够用。这套链路我前后调了三个月最大的体会是实时推荐不是把离线算法搬到 Flink 上就完事状态管理、水位线、异步查询、离线兜底每一环都有坑。先把 Kafka 的可靠性和 Flink 的状态后端配稳再调推荐算法的参数顺序反了会浪费很多时间。希望帮到你。本文还有配套的精品资源点击获取

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

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

免费获取报价 →
↑