ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink实时推荐系统实战:从数据流设计到Redis结果输出

Flink实时推荐系统实战:从数据流设计到Redis结果输出 简介基于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的binlogSink端写入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 RichParallelSourceFunctionClickLog { private volatile boolean running true; private final Random rnd new Random(); Override public void run(SourceContextClickLog 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提取的关键代码。DataStreamUserPreference prefStream source .assignTimestampsAndWatermarks( WatermarkStrategy.ClickLogforBoundedOutOfOrderness(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 ProcessWindowFunctionClickLog, UserPreference, String, TimeWindow { Override public void process(String key, Context context, IterableClickLog logs, CollectorUserPreference out) { MapInteger, Integer countMap new HashMap(); for (ClickLog log : logs) { countMap.merge(log.getCatId(), 1, Integer::sum); } ListMap.EntryInteger, Integer sorted new ArrayList(countMap.entrySet()); sorted.sort((a, b) - b.getValue() - a.getValue()); ListInteger 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”这类共现关系。流式共现统计不复杂对同一个用户把他最近点击过的商品序列保存在状态里每来一条新点击就把新商品和序列里的历史商品组成一对共现发给下游计数。DataStreamItemCoOccur coOccurStream source .keyBy(ClickLog::getUserId) .process(new CoOccurFunction()); public static class CoOccurFunction extends KeyedProcessFunctionString, ClickLog, ItemCoOccur { private transient ValueStateListLong recentItemsState; Override public void open(Configuration parameters) { ValueStateDescriptorListLong 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, CollectorItemCoOccur out) throws Exception { ListLong 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 RichSinkFunctionString { 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涨到8GFull 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对象序列化上。先用回放日志证明效果再用火焰图压性能这两步能让你的实时推荐系统少走很多弯路。希望帮到你。本文还有配套的精品资源点击获取
返回列表