ARTICLE DETAIL

资讯详情

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

Flink实时推荐系统全链路拆解:从Kafka行为接入到Redis结果落地

Flink实时推荐系统全链路拆解:从Kafka行为接入到Redis结果落地 简介一套基于Apache Flink的商品实时推荐系统完整项目源码面向大数据方向的学生、开发者和推荐系统初学者。项目以Scala为主要语言完整实现了从数据采集、预处理、特征工程到推荐算法、结果输出、系统优化的核心链路并特别提供了Flink与HBase集成读写示例以及Kafka模拟用户行为数据的生成脚本便于在本地快速搭建可运行的实时推荐Demo。压缩包共44个文件包含34个Scala源文件、SQL和HBase建表语句、Kafka数据模拟脚本、Maven工程配置等整体体积仅245KB内容紧凑却覆盖了完整工程结构。资源内部包含Flink读写HBase的模块可观察向HBase写入数据的完整过程SQL脚本提供用户表、商品表及行为表的建表语句配合Kafka模拟数据脚本能快速生成测试数据流。已有275人学习下载通过阅读源码和运行实践可以深入理解Flink DataStream API、窗口统计、状态管理以及协同过滤等推荐算法在实时场景中的落地方法是一份适合课程设计、毕业设计或项目实训的高质量参考资料。1. 基于Flink商品实时推荐系统.zip解压之后先别急着跑拿到这份“基于Flink商品实时推荐系统.zip”多数人的第一反应是解压、导入IDE、等Maven把依赖拉完然后直接点运行。我的建议恰好相反先别急着跑跑不起来的。这个压缩包里不是一套能开箱即用的服务而是一套“实时推荐链路的最小完整实现”——从Kafka里的用户行为日志到Flink实时计算用户特征和商品特征再到召回排序、写入Redis供API层读取。你需要的不是解压后的一瞬间而是先把链路图画清楚知道哪个环节缺了MySQL、Redis、Kafka或Druid缺了之后会报哪类错误。这篇笔记就按“链路设计 → 落地步骤 → 参数调整 → 排错”往下讲适合两类人一类是拿这个zip做毕设或课程设计、需要把它改造成自己业务形态的在校生另一类是刚接手实时推荐任务、想快速理解一套可用数据流的一线开发。无论哪种先管住双击运行的手我们从上到下把这条链路捋一遍。2. 实时推荐链路设计从行为埋点到召回排序的完整闭环2.1 用户行为数据从哪来Kafka 日志格式的约定Flink商品实时推荐系统处理的数据源头通常是前端或客户端上报的曝光、点击、加购、下单四类行为报文。这些报文会统一打进Kafka的topic常见的字段结构如下{ userId: u_10001, itemId: i_2048, behavior: click, scene: homepage_recommend, timestamp: 1691234567890, extra: { duration: 3200, page: detail } }这个JSON结构里有两个字段对后续实时特征计算很重要。behavior是行为类型推荐模型只应该学习正向行为但曝光数据要单独保留用于计算曝光过滤避免把用户已经看过且没点击的商品反复召回。timestamp必须带而且建议是毫秒级、事件时间语义的时间戳——后续做窗口聚合和浏览时长过滤时如果用的是数据到达Flink的processing time那么数据在kafka里积压一段时间后再消费行为顺序会失真。常见做法是把这四类行为分别路由到Kafka的不同分区甚至不同topic实时Joiner再统一做流合并。这个zip里的工程多数只有click一种行为我一般会在二次改造时把加购和下单单独拉出来因为在推荐排序阶段加购行为的权重应当远高于点击。2.2 冷热链路分开离线协同过滤与实时行为特征的接驳商品推荐系统不可能只靠实时数据。实时数据解决的是“用户刚刚看了什么、我马上给他关联推荐”的短时兴趣而离线协同过滤解决的是“和这位用户历史画像相似的人群喜欢什么”的长时兴趣。两者需要接驳接驳点是Flink任务启动时的一次全量加载。具体落地时离线部分用Spark或Hive跑一个物品协同过滤产出“相似商品对”表比如item_id, similar_item_id, score结果写入MySQL或可以直接读的HBase。实时任务启动时先用RichCoFlatMapFunction把这些相似关系加载到内存中形成一个MapString, ListSimilarItem的映射。之后每来一条用户实时行为就用当前行为商品ID去这个内存Map里捞相似商品捞到的集合作为“协同过滤召回池”再与实时热度池合并。这里有一个必须留意的边界这个内存Map不能太大。几百万对相似关系大约占几百MB堆内存可以接受上亿对就会频繁Full GC此时要换成外部缓存把相似关系放到Redis里实时任务改用异步IO查询。提示推荐系统里离线链路和实时链路的接驳最忌讳在Flink任务里实时调MySQL查相似商品。每一条用户行为都打一次MySQL行为量一上来连接池先被打爆然后Flink反压最后Kafka消费延迟以分钟计。相似关系这种静态数据启动加载或异步IO是两种可靠姿势。2.3 实时特征的计算与写入Redis里到底存什么整套链路计算出的推荐结果最终要落到API层能快速读取的存储里。Redis是这个zip项目最常见的落地存储因为推荐结果的读取延迟要求通常在10毫秒内。写入Redis的数据结构是有讲究的常见的两种设计如下。第一种是按用户维度存一个推荐列表使用Redis的String结构key为rec:user:{userId}value为JSON数组例如[ {itemId: i_3021, score: 0.92, reason: clicked_similar}, {itemId: i_1077, score: 0.87, reason: hot} ]第二种是存一个带过期时间的ZSetZADD rec:user:{userId} 0.92 i_3021。两种结构的取舍取决于API层的读法。如果一次性返回十个商品String结构更省事如果需要在Redis端做分页或截断ZSet更灵活。实际工程里我一般只写String结构因为读取逻辑简单明了且可以整体设置过期时间比如EXPIRE rec:user:{userId} 1800——推荐结果只保证半小时内有效用户再刷新一次页面时由新的实时计算重新生成。这里的核心在于实时推荐系统写入Redis的不是“模型训练出来的排序结果”而是“当前时刻结合实时行为重新算过一轮的召回排序结果”。写Redis前要把itemId去重、过滤已曝光商品、按分数降序排列这三个动作在Flink的Sink函数里完成比较合适。3. 把zip里的工程跑起来环境准备、自定义Source与Sink的落地步骤3.1 从zip解压到Flink环境搭建一次不用跑通全部的最小启动先把环境搭好。这份工程依赖的组件最少有三个Kafka、Redis、MySQL。MySQL用来存商品元数据与用户画像表Kafka用来接收模拟行为数据Redis用来存最终推荐结果。# 启动一个本地Kafka用docker compose是最省心的方式 # docker-compose-kafka.yml version: 3 services: zookeeper: image: confluentinc/cp-zookeeper:7.4.0 environment: ZOOKEEPER_CLIENT_PORT: 2181 kafka: image: confluentinc/cp-kafka:7.4.0 ports: - 9092:9092 environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1docker compose -f docker-compose-kafka.yml up -dKafka起来之后创建一个topic用于接收行为数据然后确认可以生产消费docker exec -it kafka kafka-topics --create \ --bootstrap-server localhost:9092 \ --topic user_behavior \ --partitions 3 --replication-factor 1partition数建议设成3。实时推荐任务在数据量不大时单并行度也能跑但3个分区可以让你在测试并行度调整时看到效果差异也不用因为分区太多造成Flink checkpoint过大。如果zip里自带了模拟数据发送脚本就先确认脚本往哪个topic发数据、发的数据格式是否和上文的JSON结构一致。这一步经常有偏差脚本发的字段名和Flink解析的字段名对不上运行时不会报错但所有字段都是null推荐结果永远是空列表。我一般会在这一步先用kafka-console-consumer消费几秒看看原始报文长相。3.2 自定义DataSource读取Kafka和模拟数据发生器的选择这个zip的核心必然包含一个读取Kafka的自定义DataSource。Flink官方提供的FlinkKafkaConsumer已经能满足需求但很多教学版zip喜欢写一个自定义的SourceFunction来模拟数据流原因是可以脱离Kafka独立演示。这里的风险在于演示用的自定义Source多半是无限循环生成随机JSON而不是消费Kafka里的真实数据。跑通演示容易上了生产环境还是得切回Kafka Connector。// 自定义数据源从Kafka读取用户行为JSON // DataSource与DataSink的自定义是理解Flink实时计算两条数据边界的核心 DataStreamString rawStream env.addSource( new FlinkKafkaConsumerString( user_behavior, new SimpleStringSchema(), kafkaProps ) ); // 把JSON解析成JavaBean这里用Flink自带的Jackson实现 DataStreamUserBehavior behaviorStream rawStream .map(new JsonToBehaviorFunction()) .returns(TypeInformation.of(UserBehavior.class)); // 关键参数设置事件时间与水位线保证后续窗口聚合的准确性 behaviorStream.assignTimestampsAndWatermarks( WatermarkStrategy .UserBehaviorforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) - event.getTimestamp()) );这段代码里最值得调整的是forBoundedOutOfOrderness(Duration.ofSeconds(5))这一行。5秒是“允许乱序的最大时间”如果Kafka里的行为数据时间戳偶尔乱序5秒足够覆盖绝大多数场景。调大了会导致窗口结果晚出来调小了会丢弃一部分乱序数据。流式推荐场景中晚几秒出结果通常可以接受丢数据反而会造成推荐结果空洞所以这个参数我习惯于设到10秒。3.3 自定义DataSink结果写入Redis的可靠姿势工程里自定义Sink是另一个必看模块。写入Redis的Sink需要自己实现连接池管理不能每条结果都新建连接。连接池参数直接决定高吞吐下的稳定性下面是最小可用的实现骨架public class RedisSink extends RichSinkFunctionRecommendResult { private JedisPool jedisPool; Override public void open(Configuration parameters) throws Exception { // 在open里初始化连接池而不是在构造函数或每条数据里创建 JedisPoolConfig config new JedisPoolConfig(); config.setMaxTotal(10); config.setMaxIdle(5); config.setMinIdle(2); config.setTestOnBorrow(true); // 连接池要复用否则背压一来连接先耗死 this.jedisPool new JedisPool(config, localhost, 6379, 3000); } Override public void invoke(RecommendResult value, Context context) throws Exception { try (Jedis jedis jedisPool.getResource()) { // 每个用户的推荐列表只保留一份用JSON序列化后写入 String key rec:user: value.getUserId(); String json objectMapper.writeValueAsString(value.getItems()); jedis.setex(key, 1800, json); } catch (Exception e) { // 写Redis失败不要直接抛出先检查Redis是否可用 // 这里打印日志并跳过比把整个任务搞挂更符合推荐场景 LOG.warn(write redis failed, userId{}, value.getUserId(), e); } } }Sink里的两个细节要注意。一是setex同时设置了过期时间为1800秒这比set更安全——真正生产环境里如果Flink任务挂了Redis里残留的旧推荐结果至少会自动过期不会让用户一直看到几小时前的推荐。二是testOnBorrow(true)保证每次从池里借连接时验证连接是否存活。Redis服务正常时这个参数多一点点开销但Redis重启过之后这个参数可以避免大批连接报错。提示别在Sink里把异常直接抛出去。推荐结果的Sink属于“尽力而为”的写入Redis短暂抖动导致几条写入失败用户下一次请求时重算一次就可以了。相比之下Kafka的offset提交和checkpoint稳定性更重要这两个一旦失败会造成数据重复消费或丢失。3.4 窗口聚合热门商品榜为什么能反映实时热度整个工程里最能体现“实时”二字的是热门商品榜的计算。本质上是一个滑动窗口的计数聚合Flink里用一行代码就能表达DataStreamItemHot hotStream behaviorStream .filter(behavior - behavior.getBehavior().equals(click)) .keyBy(UserBehavior::getItemId) .window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1))) .aggregate(new CountAggregate(), new HotWindowResult())这段代码里窗口的两个时间参数决定了实时热度榜的敏感度。Time.minutes(10)是窗口长度Time.minutes(1)是滑动步长含义是每1分钟计算一次过去10分钟内的点击热榜。窗口越长热度越平滑短时冲高的商品不容易立刻上榜窗口越短热度越敏感但容易出现某商品因一次小规模刷量就冲上榜的情况。商品推荐场景我一般用10分钟窗口、1分钟滑动既保证对突发热点的响应及时又不会让榜单一惊一乍。窗口聚合的结果拿到之后要和协同过滤召回、实时行为召回一起进入排序阶段。排序逻辑在推荐系统里可以很复杂但在这份zip工程里通常就是一道加权公式的ProcessFunction。把点击、加购、下单分别给不同权重再叠加热度分数最后按总分排序取TopN。这一步虽然写法简短但它决定了整个系统推荐结果的“是否像人推荐的”比任何单独模块都值得反复调参。4. 实时任务避坑指南JDBC连接器异常、Kafka积压与数据落地的典型问题4.1 JDBC连接器异常周期性写失败不是网络问题是连接池参数在裸奔把实时计算结果同时写入MySQL用于离线分析是这个zip里常见的扩展做法。但Flink JDBC连接器的报错频率在各大数据平台的热搜词里居高不下典型错误长这样Caused by: java.sql.SQLTransientConnectionException: HikariPool-1 - Connection is not available, request timed out after 30000ms这个报错的现象是任务运行前半小时正常之后周期性出现写失败而且每次失败间隔和窗口触发周期高度重合。原因基本是连接池的最大连接数配置小于并发写入的Task数量窗口一批数据到来时多个并行子任务同时申请连接池里的连接被借完新请求排队直到超时。解决方法是调整JDBC连接器的参数-- 在Flink SQL中使用JDBC连接器时通过SQL Hints调整连接池参数 INSERT INTO mysql_analytics_table /* OPTIONS(sink.buffer-flush.max-rows 200, sink.buffer-flush.interval 5s, sink.max-retries 3) */ SELECT * FROM realtime_features;这里的三个参数各解决一个问题。sink.buffer-flush.max-rows控制攒够多少行才刷一次默认是1000但MySQL服务端如果性能一般攒太多行一次写入容易造成锁等待调低到200可以缓解。sink.buffer-flush.interval控制最多等多久必须刷一次设5秒保证延迟可控。sink.max-retries是写入失败后的重试次数默认值偏小数据库抖动时容易直接把checkpoint搞失败。注意如果改完参数仍然周期性报错去查MySQL的max_connections和wait_timeout。很多所谓“JDBC连接器异常”根源是数据库侧的wait_timeout默认8小时连接池里的连接长时间空闲被MySQL服务端断开连接池却还认为连接存活。这属于经典的“两端参数不匹配”排查方向从一开始就要把范围扩大到数据库服务端。4.2 Kafka消息积压反压源头往往不在Source而在下游Sink实时推荐任务有一个高频现象Kafka消费延迟从秒级涨到分钟级kafka-consumer-groups查看lag持续上涨。新手第一反应是加大并行度于是把Source的并行度从3改到12但很快发现lag不降反升。真实原因通常是下游有慢操作要么是排序阶段里有外部调用要么是Sink写入Redis时单条同步写太慢。用Flink的Web UI排查时看每个算子的BackPressure指标。如果Sink: RedisSink显示HighSource反倒正常说明背压是从尾部往回传导的。此时解决方案不是加Source并行度而是给Sink加批量写入能力// 解决方案把逐条写Redis改成先攒一批再批量写 public class BatchRedisSink extends RichSinkFunctionRecommendResult { private transient ListRecommendResult buffer; private static final int BATCH_SIZE 100; private static final long BATCH_INTERVAL_MS 2000; Override public void open(Configuration parameters) { this.buffer new ArrayList(); } Override public void invoke(RecommendResult value, Context context) { buffer.add(value); if (buffer.size() BATCH_SIZE) { flush(); } } private void flush() { // 批量写Redis用pipeline减少RTT try (Jedis jedis jedisPool.getResource()) { Pipeline pipeline jedis.pipelined(); for (RecommendResult result : buffer) { pipeline.setex(rec:user: result.getUserId(), 1800, JsonUtils.toJson(result.getItems())); } pipeline.sync(); } buffer.clear(); } }这个改动的关键点是pipeline。原本100条结果要100次RTTpipeline可以合并成一次网络往返吞吐量能提升一个数量级。代价是延迟从单条立即写入变成最多2秒的攒批窗口对推荐结果而言完全可接受。做这类优化时心里要有一杆秤实时推荐系统的实时性指的是行为发生后几秒内能反映到结果里而不是每一条计算结果都即刻可查。4.3 Sink Hive表数据不入表分区提交是个假象这个zip如果扩展了Hive数仓链路会遇到一个很典型的“数据不入表”现象flink任务在Web UI上显示sink成功checkpoint也正常但到Hive分区目录里看.staging文件一堆正式分区数据却是空的。这是把实时数据写入Hive表时的经典坑原因基本可以锁定在“流式写入分区文件需要streaming阶段自动提交”。排查顺序如下。先看表属性建表时是否声明了streaming相关参数。Hive表需要开启文件提交机制否则Flink的FileSink会一直写.staging临时文件永远不会rename成正式文件-- 建表时必须指定streming相关属性缺了这个文件不会从staging转正 CREATE TABLE recommendation_log ( user_id STRING, item_id STRING, recommend_type STRING, ts BIGINT ) PARTITIONED BY (dt STRING, hh STRING) STORED AS PARQUET TBLPROPERTIES ( streaming true, auto-compaction false );再看Flink SQL或者DataStream写Hive时是否设置了分区提交的触发策略。常见做法是// DataStream API写Hive时开启分区提交的两种触发方式 // 1. 基于处理时间每10分钟提交一次 // 2. 基于checkpoint每个checkpoint尝试提交 FileSinkString sink FileSink.forBulkFormat( path, new ParquetRowDataBuilder(...) ) .withPartitionCommitter(new HivePartitionCommitter(conf, catalogTable)) .withPartitionCommitterStrategy(new MetastoreCommitPolicy()) .build();很多教学工程只实现了写入逻辑没有实现PartitionCommitter。这种情况下数据确实写进了Hive的临时目录Web UI也显示写出去了但外表看不到任何数据。解决方案就是在FileSink上补上分区提交策略同时把checkpoint间隔设置为与分区提交周期匹配的时长。如果不想动代码最土的办法是定期用MSCK REPAIR TABLE修复分区但不推荐上生产治标不治本。5. 用CEP实现“几分钟内浏览过A又浏览过B”的关联推荐一个值得深度改造的进阶方向5.1 为什么CEP比窗口聚合更适合商品关联推荐协同过滤做的是“看了A的人还看了B”实时热度做的是“现在大家都在看什么”但这两者都缺一种能力实时捕获“这位用户刚刚在短时间内连续浏览了哪些商品把这些商品关联起来”。比如用户两分钟内依次查看了相机、镜头、三脚架这时候最合理的推荐是把这三者的配件推荐出来。用滑动窗口聚合来处理这个场景非常别扭因为窗口边界是死的而CEP可以定义“A出现后10分钟内出现B”这类事件模式。Flink CEP的代码能直观表达这种业务规则。5.2 一条可运行的CEP关联规则// 定义模式同一用户在10分钟内浏览了商品A又浏览了商品B PatternUserBehavior, UserBehavior pattern Pattern .UserBehaviorbegin(first) .where(new SimpleConditionUserBehavior() { Override public boolean filter(UserBehavior behavior) { return click.equals(behavior.getBehavior()); } }) .next(second) .where(new SimpleConditionUserBehavior() { Override public boolean filter(UserBehavior behavior) { return click.equals(behavior.getBehavior()); } }) .within(Time.minutes(10)); // 在click流上应用这个模式输出一次关联事件 DataStreamString matchedStream CEP.pattern(behaviorStream, pattern) .inProcessingTime() .select((MapString, ListUserBehavior patternMap) - { UserBehavior first patternMap.get(first).get(0); UserBehavior second patternMap.get(second).get(0); return first.getItemId() - second.getItemId(); });这段代码里的within(Time.minutes(10))是“关联间隔”的核心参数。设太大会把用户半小时前看过的商品和现在浏览的商品强行关联关联噪音高设太小又捕捉不到真实的连续性浏览行为。我通常在电商场景先用5到10分钟起步然后观察推荐结果里的关联商品点击率来反向调整。inProcessingTime()代表用处理时间做模式匹配实时性高但结果不可精确重放如果后续要做效果对比实验换成inEventTime()配水位线更严谨。这块改造成本不高但收益很明显推荐结果里会出现一批“基于用户当前浏览序列”的动态关联商品。协同过滤提供的相似商品是相对静态的而CEP关联规则让每个用户看到的是“跟着本次浏览路径走的推荐”这是实时推荐最容易体现差异化的能力。玩到这一步这个zip工程已经不只是“能跑起来”而是一个可以往生产形态演变的骨架了。回头看我自己的实践经历每次做实时推荐改造最深的教训都是同一个不要把实时推荐当成一个纯计算问题它是一个工程链路问题Kafka里数据的质量、Redis里结果的过期策略、MySQL连接池的脾气任何一个环节拉胯Flink计算再正确也白搭。看不懂的线上故障十有八九出在链路两端而不是计算引擎本身。希望帮到你。本文还有配套的精品资源点击获取
返回列表