☰
Flink实时推荐系统实战:从数据流设计到Redis结果输出
2026/10/3 3:19:04 网站建设 项目流程

简介:基于Flink实现的商品实时推荐系统,面向大数据开发工程师与推荐系统学习者,用于解决实时商品热度统计、日志分析与个性化推荐等问题。系统以Flink为处理核心,实时统计商品热度并写入Redis缓存;同时采集分析用户日志,将画像标签与实时记录存入HBase。用户发起推荐请求时,系统依据用户画像对热度榜重新排序,并结合协同过滤与标签推荐两个模块,为榜单中每个商品补充关联产品,形成更完整的推荐列表。资源包为zip格式,共109个文件,以68个Java源码为主体,覆盖推荐服务实现、用户评分服务等核心逻辑;另含SQL建表脚本、HBase建表语句、Kafka模拟数据生成脚本、前端展示页面、Spring配置与说明文档,便于从数据接入到结果输出的全链路理解。压缩包仅3.74MB,轻量且目录清晰。目前已有130人学习下载,适合具备一定Flink基础、希望获得完整可运行方案并快速上手实时推荐系统开发的读者。

1. 实时推荐系统:不要等到离线跑完才想起用户已经走了

用户点开一个商品详情页后,推荐栏要在几百毫秒内给出“接下来买什么”,这个需求离线推荐很难满足。离线任务凌晨跑一次,中午才出结果,用户下午看到的还是昨天凌晨的兴趣,哪怕他刚点击了一条新商品,推荐列表也纹丝不动。基于Flink实现的商品实时推荐系统,正是把“从点击到推荐”这个过程压缩到秒级:Flink吃掉Kafka里的行为日志,实时维护每个用户的兴趣状态,再做召回、排序,最终把TopN结果写回Redis。它适合正在被离线延迟困扰、又不想一上来就上大模型排序的团队,也适合想用真实业务场景练手Flink的开发者。下面这条链路是我在多个项目里反复用过的,照着搭能把上线时间缩短一大半。

2. 先想清楚拓扑再写代码:实时推荐系统的数据流与选型

写实时推荐系统最容易犯的错,是刚学会Flink API就急着在IDE里写逻辑,结果数据源、状态、存储没有一个对得上。我的习惯是先画一条数据流图,把每一条数据在链路里扮演什么角色想清楚,再开始写代码。这一章会从数据接入一直讲到结果存储,顺便把选型理由讲透。

2.1 一条用户行为从点击到推荐结果要走完哪几站

一条原始的点击行为,通常从前端埋点进入消息队列Kafka。Kafka在这里承担两个职责:一是削峰填谷,晚高峰的流量是平时的几十倍,Flink直接接数据库会被压垮;二是保存一份可重放的行为历史,Flink任务升级或者从checkpoint恢复时,可以重新消费。Flink从Kafka的某个topic读取日志,完成清洗、特征计算、召回、排序,最后把每个用户的TopN商品列表写入Redis。推荐服务收到页面请求时只是从Redis按用户ID取值,不需要自己算,RT能稳定在几十毫秒。

为什么不直接用Spark Streaming?微批模型天然有秒级延迟,而“用户刚点击了什么”这类信号的有效期可能只有几分钟,延迟5秒和延迟500毫秒,体验差距很大。Flink的DataStream API和事件时间处理更贴合这个场景。另外,实时推荐的很多算子需要保存用户维度状态,比如“最近点击的20个商品”“最近下单的类目”,Flink把状态管理做进了核心算子,代码写起来比Spark更顺手。如果你只是做离线特征回溯,Spark没问题;但要做实时更新,我建议从Flink开始。

2.2 离线与实时并行:用Flink CDC同步订单数据

如果只用点击日志做推荐,离业务需求还差得远。用户下单是最强的正向信号,可订单数据通常躺在MySQL里。常见做法是用Flink CDC把订单库的变更实时同步到Kafka,再让推荐任务订阅Kafka,这样实时推荐系统就能同时感知点击、加购和订单行为。Flink CDC的部署也被称为“Flink CDC Pipeline”,简单说就是Source端解析MySQL的binlog,Sink端写入Kafka或数据湖,中间不需要业务方改一行代码。

下面是一段常见的MySQL CDC建表语句,启动后Flink会自动读取订单表的初始数据,并持续跟踪后续变更:

CREATE TABLE orders ( order_id BIGINT PRIMARY KEY NOT ENFORCED, user_id BIGINT, item_id BIGINT, status STRING, amount DECIMAL(10,2), create_time TIMESTAMP(3), WATERMARK FOR create_time AS create_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'mysql-server', 'port' = '3306', 'username' = 'flink_user', 'password' = '***', 'database-name' = 'shop', 'table-name' = 'orders', 'scan.startup.mode' = 'initial' );

逻辑说明:这段SQL定义了一张动态表,Flink CDC会按binlog顺序把orders表的全部变更变成流。scan.startup.mode填initial,表示任务首次启动先做一次全表快照,之后增量读binlog,适合订单表数据量不大、需要完整历史的情况;如果只关心新增变更,可以改成latest-offset。注意字段类型要和MySQL对齐,比如DECIMAL对应DECIMAL(10,2),如果写成DOUBLE,反序列化和水位线计算都会出问题。CDC任务在快照阶段会对源库产生压力,建议在从库上跑或者选业务低峰期启动。

2.3 环境准备:Flink安装配置到部署的捷径与两个参数

很多团队在“Flink安装配置到部署”这一步就卡住了,其实本地验证根本不需要搭生产集群。开发阶段用standalone集群,把Flink解压后改一个yaml就能跑。最值得调的两个内存参数是jobmanager.memory.process.size和taskmanager.memory.process.size。JobManager只负责调度,不需要给太多;TaskManager才是执行任务的地方,要按数据量和状态大小来给。

我常用的开发配置如下:

jobmanager.memory.process.size: 4096m taskmanager.memory.process.size: 8192m taskmanager.numberOfTaskSlots: 4 parallelism.default: 2

参数说明:taskmanager.numberOfTaskSlots决定一个TaskManager里能放多少个子任务,建议和机器CPU核数一致,不要盲目设大。parallelism.default是全局默认并行度,推荐任务通常要和Kafka分区数对齐,比如Kafka topic有6个分区,Source并行度就设6,否则会有分区数据空闲。jobmanager.memory.process.size我一般不超过8G,因为它不负责状态存储。这套配置在本地足够跑通测试;生产用Flink on YARN时,把并行度提上去,并把RocksDB的托管内存比例调大。

选型确定后,链路里还缺两个关键部分:Flink怎么拿到行为数据、怎么把计算结果输出到Redis。下一章直接用代码实现这两个环节。

3. 用Flink把用户兴趣算出来:自定义DataSource与特征计算

很多教程一上来就让你连Kafka,但本地没有Kafka环境时,最方便的做法是先用自定义Data Source造一段模拟点击流,把主链路跑通,再切换成真实Kafka源。这也是“Flink实现自定义Data Source”这个进阶点的实际用处:当官方connector覆盖不了你的数据格式时,你需要自己写Source。

3.1 先造数据:自定义Data Source模拟用户点击流

Flink允许直接实现SourceFunction或继承RichParallelSourceFunction来生成数据。下面这段代码每100毫秒生成一条点击日志,包括用户ID、商品ID、商品类别ID、行为类型、事件时间戳。

public class ClickSource extends RichParallelSourceFunction<ClickLog> { private volatile boolean running = true; private final Random rnd = new Random(); @Override public void run(SourceContext<ClickLog> ctx) throws Exception { String[] users = {"u_1001", "u_1002", "u_1003", "u_1004", "u_1005"}; int[] items = {101, 202, 303, 404, 505}; int[] cats = {1, 2, 3, 4, 5}; while (running) { long ts = System.currentTimeMillis(); ctx.collect(new ClickLog( users[rnd.nextInt(users.length)], items[rnd.nextInt(items.length)], cats[rnd.nextInt(cats.length)], "click", ts )); Thread.sleep(100); } } @Override public void cancel() { running = false; } }

逻辑说明:这里继承的是RichParallelSourceFunction,可以并行读取,并行度由下游算子决定。running用volatile修饰,调用cancel()时停止循环。事件时间直接用当前毫秒,省去水印生成;真实项目里应该从Kafka消息里解析业务时间戳。参数上,Thread.sleep(100)控制发射频率,想观察水位线和迟到数据,可以把ts改成System.currentTimeMillis() - rnd.nextInt(5000),让数据乱序,再在流上配WatermarkStrategy。

3.2 用事件时间窗口统计用户类别偏好

模拟数据出来后,第一步是按用户维度统计“最近15分钟点击最多的三个商品类别”。这里我直接使用事件时间滚动窗口,相比在每个process里手工维护定时器更标准,而且窗口天然支持迟到数据。下面是窗口处理和TopN提取的关键代码。

DataStream<UserPreference> prefStream = source .assignTimestampsAndWatermarks( WatermarkStrategy.<ClickLog>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((log, ts) -> log.getTs()) ) .keyBy(ClickLog::getUserId) .window(SlidingEventTimeWindows.of(Time.minutes(15), Time.minutes(5))) .process(new TopCategoryWindowFunction()); public static class TopCategoryWindowFunction extends ProcessWindowFunction<ClickLog, UserPreference, String, TimeWindow> { @Override public void process(String key, Context context, Iterable<ClickLog> logs, Collector<UserPreference> out) { Map<Integer, Integer> countMap = new HashMap<>(); for (ClickLog log : logs) { countMap.merge(log.getCatId(), 1, Integer::sum); } List<Map.Entry<Integer, Integer>> sorted = new ArrayList<>(countMap.entrySet()); sorted.sort((a, b) -> b.getValue() - a.getValue()); List<Integer> topCats = sorted.stream() .limit(3) .map(Map.Entry::getKey) .collect(Collectors.toList()); out.collect(new UserPreference(key, topCats, context.window().getEnd())); } }

逻辑说明:SlidingEventTimeWindows.of(Time.minutes(15), Time.minutes(5))表示每5分钟输出过去15分钟的用户偏好,既不会太滞后,又能捕捉短期兴趣变化。窗口结束时间作为结果的时间戳,方便下游判断数据新旧。这里要注意ProcessWindowFunction会把窗口内所有数据暂存在内存里,如果用户量极大,建议先在AggregateFunction里做增量计数,再用ProcessWindowFunction输出TopN,避免单窗口数据量过大。amforBoundedOutOfOrderness(Duration.ofSeconds(5))这个5秒表示容忍迟到5秒,超过的水印会直接丢弃,业务上可接受的延迟阈值要自己测。

3.3 实时召回:用共现关系生成候选商品

类别偏好是召回的上层过滤器,真正能带来点击的往往是“看了A的人也会看B”这类共现关系。流式共现统计不复杂:对同一个用户,把他最近点击过的商品序列保存在状态里,每来一条新点击,就把新商品和序列里的历史商品组成一对共现,发给下游计数。

DataStream<ItemCoOccur> coOccurStream = source .keyBy(ClickLog::getUserId) .process(new CoOccurFunction()); public static class CoOccurFunction extends KeyedProcessFunction<String, ClickLog, ItemCoOccur> { private transient ValueState<List<Long>> recentItemsState; @Override public void open(Configuration parameters) { ValueStateDescriptor<List<Long>> desc = new ValueStateDescriptor<>( "recentItems", new ListTypeInfo<>(Types.LONG)); StateTtlConfig ttl = StateTtlConfig.newBuilder(Duration.ofHours(1)) .setUpdateType(StateTtlConfig.UpdateType.OnReadAndWrite) .cleanupIncrementally(1000, true) .build(); desc.enableTimeToLive(ttl); recentItemsState = getRuntimeContext().getState(desc); } @Override public void processElement(ClickLog log, Context ctx, Collector<ItemCoOccur> out) throws Exception { List<Long> recent = recentItemsState.value(); if (recent == null) recent = new ArrayList<>(); long currentItem = log.getItemId(); for (Long before : recent) { if (!before.equals(currentItem)) { out.collect(new ItemCoOccur(before, currentItem, 1L)); } } recent.add(currentItem); if (recent.size() > 20) { recent.remove(0); } recentItemsState.update(recent); } }

逻辑说明:recentItemsState保存的是每个用户最近20个点击过的商品,状态TTL设为1小时,超过1小时自动清理。每来一次新点击,就和历史商品组成(before, current)对,下游再keyBy(before).sum(1)就能得到“商品A后接了商品B”的共现次数。这里有一个容易被忽略的点:状态里只保留20个商品,是为了控制状态大小;如果用户行为很长,旧商品的共现信息会自然流失,但流式推荐要的就是短期兴趣,所以别把状态开太大。

上述代码跑通后,你已经有了“用户最近喜欢的类目”和“商品与商品之间的关系”。下一步就是排序和输出,把候选集变成一个用户真正会看的列表。

4. 排序与结果写出:从候选集到Redis只有几步

召回只是把可能感兴趣的商品捞回来,用户最终看到什么还取决于排序。实时推荐初期不建议直接上模型排序,先用一套可解释的规则打底,等数据积累了再替换成模型。这一章讲清楚规则怎么设计、结果怎么安全地写进Redis。

4.1 规则排序:为什么先用规则而不是模型

在数据量和业务复杂度不高时,规则排序比模型排序更快见效,也更容易让运营理解“为什么推荐了A没推荐B”。一套常用的实时排序规则是:候选商品必须落在用户最近偏好的3个类目里;过滤掉用户最近7天下过单的商品;再按加权分数排序。

我用的打分公式如下:

score = 0.4 * hotScore(item) + 0.3 * coOccurScore(item, userRecentItems) + 0.2 * categoryMatchScore(item, userPrefCats) + 0.1 * freshScore(item)

四个分数含义很直白:hotScore是商品近1小时热度,用点击量归一化;coOccurScore是当前用户最近点击商品与候选商品的总共现值;categoryMatchScore是候选商品类目和用户偏好类目的交集数;freshScore是商品上架时长对分数的衰减。权重先用经验值,上线后看推荐位点击率再调。用DataStream实现时,把召回结果流转成ScoreItem,然后用keyBy(userId).process做TopN排序,最后输出JSON字符串。这样就算之后换模型排序,也只是把这个算子内部替换掉,上下游都不用动。

4.2 把结果写入Redis:自定义Sink与TTL的坑

实时推荐结果要快,Redis是首选。我习惯把每个用户的Top50商品列表写成一条JSON字符串,key为reco:user:{userId},TTL设为10分钟。为什么必须设TTL?用户兴趣每几分钟就会变,旧结果如果一直留着,会在Redis里堆积成脏数据,业务侧读到的永远不是最新兴趣。

下面是一个简单的Redis Sink示例:

public class RedisRecommendSink extends RichSinkFunction<String> { private transient JedisPool pool; private final int ttlSeconds; public RedisRecommendSink(int ttlSeconds) { this.ttlSeconds = ttlSeconds; } @Override public void open(Configuration parameters) { JedisPoolConfig config = new JedisPoolConfig(); config.setMaxTotal(20); config.setMaxIdle(10); config.setMinIdle(5); pool = new JedisPool(config, "redis-host", 6379, 3000, "password", 0); } @Override public void invoke(String json, Context ctx) { String userId = extractUserIdFromJson(json); try (Jedis jedis = pool.getResource()) { jedis.setex("reco:user:" + userId, ttlSeconds, json); } catch (Exception e) { // 记录日志后继续,别让下游异常阻塞主线程 } } @Override public void close() { if (pool != null) pool.close(); } }

逻辑说明:这个Sink在open里初始化连接池,避免每条数据新建Jedis连接。setex会用原子操作设置值并带上过期时间,比先set再expire更安全。extractUserIdFromJson需要你用Fastjson或Jackson解析,实际项目中也可以用Redis的Hash结构按字段存,但字符串JSON在推荐服务侧解析最方便。ttlSeconds一般设600,如果业务希望用户重新打开App就刷新推荐,可以缩短到300秒。

4.3 用Flink SQL做维表关联:JDBC连接器常见配置

推荐排序需要商品标题、价格、品牌等静态属性,这些信息存在MySQL。如果在DataStream的map里逐条同步查库,连接数会瞬间打满。常见做法是用Flink SQL的维表JOIN,让Flink自己维护Lookup缓存。下面这张商品维表是我常用的配置:

CREATE TABLE product_dim ( item_id BIGINT PRIMARY KEY, category_id INT, title STRING, price DECIMAL(10,2), brand STRING ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://mysql-server:3306/shop', 'table-name' = 'product_dim', 'username' = 'flink_user', 'password' = '***', 'lookup.cache.max-rows' = '10000', 'lookup.cache.ttl' = '10 min' );

然后在推荐SQL里用FOR SYSTEM_TIME AS OF去关联这张表。lookup.cache.max-rows和lookup.cache.ttl是最需要调的两个参数:max-rows控制缓存商品条数,ttl控制缓存过期时间。很多所谓的“Flink的JDBC连接器异常”,其实是缓存太小导致每次查询都回源MySQL,连接池被耗尽。生产上建议先把lookup.cache.ttl设为10分钟,观察数据库负载再继续调。另外,Flink SQL的JDBC连接器会自动使用连接池,你不需要在代码里手动建连接,但要记得给这张表数据库账号开通只读权限。

到这里,一个实时推荐主链路已经闭环:行为日志进Flink,算出偏好和召回,排序后写Redis,推荐服务读表。但上线前必须正视几个高频坑,下面这5个是真实环境里最容易遇到的。

5. 避坑/常见问题/排查:实时推荐系统上线前一定要看的5个坑

5.1 状态无限增长,任务跑了三天就OOM

现象:TaskManager堆内存从2G涨到8G,Full GC频繁,最后Container被YARN杀掉,任务一直重启。

原因:用户偏好或最近点击商品的状态没有设置TTL。Flink的Keyed State默认永久保留,只要key一直有新数据,状态就会无限增长。

解决:给所有用户维度的State都配上TTL。示例如下:

StateTtlConfig ttl = StateTtlConfig.newBuilder(Duration.ofHours(1)) .setUpdateType(StateTtlConfig.UpdateType.OnReadAndWrite) .cleanupIncrementally(1000, true) .build();

cleanupIncrementally(1000, true)表示每处理1000条数据检查一次清理,true代表后台线程也会触发。如果你用RocksDB存储状态,还可以调state.backend.rocksdb.ttl.compaction.filter.enabled,让底层压缩时物理删除过期数据。记住,只要状态key是用户ID或商品ID,就必须有TTL,否则推荐系统跑得越久越卡。

5.2 热门商品造成数据倾斜,推荐结果卡顿

现象:某个算子的并行度是6,但只有一个subtask负载特别高,背压面板上它的水位线明显落后,其他5个subtask空闲。

原因:热门爆款商品的点击量远大于普通商品,Flink按键分组时,同一个itemId的所有数据都进了同一个子任务。

解决:对需要聚合的itemId做“加盐”处理。比如在itemId后面拼一个随机后缀0~N,先分散聚合,再合并去掉后缀。我在实时共现统计里常用这个办法:先keyBy(itemId + "_" + rnd.nextInt(8))做局部计数,再keyBy(itemId)汇总。注意加盐会改变数据语义,只适用于可交换的统计操作;如果业务上必须保证同一用户的状态一致,盐要按userId维度做,不能按商品维度。

5.3 Flink Sink Hive表数据不入表,任务成功率却100%

现象:实时推荐任务运行正常,Checkpoint成功率100%,但去Hive分区表查询,一条数据都没有。

原因:Flink写Hive用的是文件流式写入,数据先写到临时目录,只有checkpoint完成时才会把pending文件转为正式文件。如果你checkpoint间隔设成几分钟,或者分区提交逻辑没配对,数据就一直留在临时目录等待提交,表里自然查不到。

解决:先把execution.checkpointing.interval调成60秒,让文件尽快提交。检查Hive表的sink参数,如果用的是Flink SQL写Hive,需要在建表时指定'sink.partition-commit.trigger' = 'process-time',并给'sink.partition-commit.delay'设一个正数。如果还在用DataStream的StreamingFileSink,要确认withOutputFormat和withBucketAssigner是否正确。最简单的方法:先写到一个非分区表,确认数据能落,再上分区表,避免排错时把分区和文件提交两件事混在一起。

5.4 JDBC连接器异常:连接池不够、连接被回收

现象:维表JOIN时偶发“Connection is not available, request timed out after 30000ms”,服务间歇性报错,过一会自己恢复。

原因:JDBC连接器底层连接池默认值太小,高并发维表查询把连接占满,新的查询只能等超时。另外,如果MySQL的wait_timeout设置得短,空闲连接会被服务端断开,连接池里残留的坏连接也会导致异常。

解决:调大lookup.cache.max-rows和lookup.cache.ttl,减少回源次数;同时调大连接池上限,把maximum-pool-size从默认10改成50,并设置connection-timeout为3秒。Flink SQL的JDBC连接器有些版本不暴露连接池配置,这时可以退回到DataStream API的JdbcLookupFunction,自己控制连接池。还有一个容易被忽略的点:维表数据量不大时,干脆把整表加载到Flink的广播状态里,完全绕开连接池,性能最好。

5.5 重启后推荐结果重复或丢失

现象:从checkpoint恢复后,用户看到的推荐列表还是旧数据,甚至同一条结果出现两次。

原因:Redis结果集只做了setex覆盖,没有清理上一轮结果;更麻烦的是Flink恢复后读取Kafka offset时有重复消费,而行为消息里没有唯一的业务主键,下游也就无法去重。

解决:在写Redis的JSON里带上一个batchId或generate_time字段,推荐服务读取时只认最新的批次;或者给结果key加一个版本号,比如reco:user:{userId}:{batchId},写完新数据后再删旧key。如果你用的是Kafka的结果topic,建议把topic格式配成upsert,以userId为主键,Flink能保证最后一条数据覆盖前一条,从源头降低重复。

6. 进阶:如何验证实时推荐效果并用火焰图定位背压

6.1 用词频统计验证环境,再跑推荐主链路

刚搭好Flink环境时,别直接跑推荐任务。先用一个最简单的“Flink实时计算-词频统计初体验”验证集群:在终端执行nc -lk 9999,然后提交WordCount任务往9999端口发字符串,看输出是否正常。这个过程十分钟内能验证安装配置、提交命令、Web UI日志是否正常。之后再把第三章的ClickSource换成Kafka源,按keyBy(userId).process跑推荐主链路,确认能从Kafka持续消费并写出Redis。

6.2 用Flink火焰图定位背压和瓶颈

推荐任务最常见的故障是“任务没报错,但Kafka积压越来越多”。先在Web UI看Backpressure状态,如果某一层是HIGH,就需要看火焰图。用async-profiler挂到TaskManager JVM上,采样一分钟左右,就能看到真正的CPU热点。我遇到最多的是这两类:TypeSerializer相关方法耗时高,说明POJO序列化是瓶颈,解决办法是改用Avro或自定义序列化器;RocksDB状态访问耗时高,说明状态太大或读写太频繁,可以调大托管内存、减少状态字段。火焰图是用来“看事实”的,不要凭感觉调参。

6.3 推荐效果验证:回放历史日志

最后一步,验证推荐系统是不是真的有效。最简单的办法是回放历史日志:把某一天的点击行为从Kafka源头重新推进Flink任务,记录每个用户当时生成的推荐列表,再和真实App日志比对,看用户有没有点推荐位商品。有了这个回放数据,就能算推荐位曝光点击率和转化率。我做过一个实时推荐任务,上线前只看延迟没看状态清理,三天后OOM;后来靠火焰图发现一半CPU花在Java对象序列化上。先用回放日志证明效果,再用火焰图压性能,这两步能让你的实时推荐系统少走很多弯路。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询