尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

基于Flink的全端用户画像实时推荐系统:从数据接入到在线查询的完整实践

基于Flink的全端用户画像实时推荐系统:从数据接入到在线查询的完整实践 简介这份资源是《基于Flink全端用户画像商品推荐系统》的完整项目源码包面向学习大数据实时处理与推荐算法的计算机专业学生及开发者可用于课程设计、毕业设计或技术练手。项目以Apache Flink为核心引擎串联数据采集、实时清洗聚合、动态用户画像构建与商品推荐展示覆盖协同过滤、矩阵分解等推荐思路帮助理解流批一体架构在电商场景中的落地方式。压缩包共27个文件以25个Java源码和2个XML配置为主Java文件承载数据流处理与推荐逻辑XML负责Maven依赖与工程构建整体约24KB结构轻量便于快速导入IDE运行调试。目前已有153人学习下载。通过研读源码读者可掌握Flink DataStream API的使用、用户行为数据的实时计算链路、画像更新与推荐策略的动态调整方法并借鉴其模块划分与工程组织方式为后续搭建实时推荐系统提供可复用的参考骨架。1. 从「标签堆叠」到「实时意图」Flink 全端用户画像推荐到底在解决什么很多团队做推荐系统的起点是一张离线的用户标签宽表昨天算好的性别、年龄、偏好类目今天原封不动地喂给召回和排序。结果就是用户上午刚搜过婴儿车下午首页还在推机械键盘——因为画像更新链路是 T1 的推荐自然慢半拍。基于 Flink 全端用户画像商品推荐系统要解决的正是这条链路的时效问题把埋点、订单、加购、搜索这些全端行为用 Flink 做实时聚合沉淀成可被在线服务直接查询的画像再驱动商品推荐。它适合两类人一类是手里已经有离线画像、想把它升级成实时链路的推荐工程同学另一类是刚接触 Flink想找一个能串起「数据接入—状态计算—结果存储—在线查询」完整闭环的练手项目。读完你应该能自己搭出一条最小可跑的实时画像管道并知道哪些参数一改就翻车。2. 全端画像的数据链路怎么拆从埋点到推荐结果的四段式2.1 为什么是「全端」而不是「单端」单端画像只看 App 或只看小程序用户换个端行为就断了。全端的意思是同一用户在 App、H5、小程序、PC 上的行为要归一到同一个user_id下。常见做法是埋点里带device_id和登录后的user_id用 Flink 做一次 ID-Mapping 的流式关联把匿名期的行为和登录后的账号合并。这一步不做后面所有画像都是残缺的。链路我一般拆成四段接入层Kafka 收埋点和业务 binlog、计算层Flink 做窗口聚合和标签计算、存储层HBase/Redis/ClickHouse 存画像、服务层SpringBoot 暴露查询接口给推荐引擎。四段里最容易出问题的是计算层和存储层的衔接后面会重点讲。2.2 接入层Kafka Topic 怎么规划埋点数据建议按端拆 Topic比如ods_behavior_app、ods_behavior_h5业务数据用 Flink CDC 从 MySQL binlog 抓。Topic 分区数按峰值 QPS 估单分区 510MB/s 是常见经验值。分区键用user_id保证同一用户行为进同一分区后续 keyBy 不会乱序。# 建埋点 Topic6 分区副本 2 kafka-topics.sh --create \ --bootstrap-server localhost:9092 \ --topic ods_behavior_app \ --partitions 6 \ --replication-factor 2 # 业务库 binlog TopicFlink CDC 会写入 kafka-topics.sh --create \ --bootstrap-server localhost:9092 \ --topic ods_order_binlog \ --partitions 3 \ --replication-factor 2分区数不是越多越好分区过多会让 Flink 的 checkpoint 变慢状态后端压力大。副本数生产环境至少 2测试环境 1 即可。2.3 计算层Flink 作业的骨架计算层核心是两件事行为聚合和标签计算。行为聚合用滚动窗口或滑动窗口统计近 1 小时、近 24 小时的点击/加购/下单次数标签计算基于聚合结果打标签比如「近 1 小时点击母婴类目 5 次」就标记为「母婴高意向」。// Flink 1.17 作业骨架行为聚合 标签输出 StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(60000); // 1 分钟一次 checkpoint env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); DataStreamBehaviorEvent behavior env .addSource(new FlinkKafkaConsumer(ods_behavior_app, new BehaviorSchema(), props)) .assignTimestampsAndWatermarks( WatermarkStrategy.BehaviorEventforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((e, ts) - e.getEventTime()) ); // 按用户开 1 小时滑动窗口每 5 分钟滑动一次 DataStreamUserTag tags behavior .keyBy(BehaviorEvent::getUserId) .window(SlidingEventTimeWindows.of(Time.hours(1), Time.minutes(5))) .aggregate(new BehaviorAggregator(), new TagProcessWindowFunction()); tags.addSink(new HBaseSink()); // 画像写入 HBaseenableCheckpointing(60000)是 1 分钟一次生产环境可以调到 35 分钟减少对存储的压力。forBoundedOutOfOrderness(Duration.ofSeconds(5))是允许 5 秒乱序埋点延迟高的场景可以放宽到 30 秒但窗口结果会晚出。SlidingEventTimeWindows用事件时间别用处理时间否则重跑数据结果不一致。2.4 存储层画像存哪、怎么查画像存储选型看查询模式点查按 user_id 查标签用 HBase 或 Redis范围/聚合查询统计某标签人群规模用 ClickHouse。我一般 HBase 存明细标签Redis 存热点用户的画像缓存ClickHouse 做离线分析。写入 HBase 时 RowKey 设计成user_id 标签类型反序避免热点。3. 用 Flink 把 MySQL 画像维表同步到 ClickHouse 的实操3.1 为什么画像维表要同步到 ClickHouse用户画像里有一部分是慢变维表比如用户注册信息、会员等级这些存在 MySQL 里。推荐引擎做特征拼接时需要快速关联MySQL 扛不住高并发点查所以常见做法是用 Flink CDC 把 MySQL 同步到 ClickHouse利用 ClickHouse 的列存做快速过滤和聚合。3.2 Flink CDC 同步作业的写法-- Flink SQL 方式MySQL CDC 到 ClickHouse CREATE TABLE mysql_user_profile ( user_id BIGINT, member_level INT, register_time TIMESTAMP(3), PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname mysql-host, port 3306, username cdc_user, password cdc_pass, database-name profile_db, table-name user_profile ); CREATE TABLE clickhouse_user_profile ( user_id BIGINT, member_level INT, register_time TIMESTAMP(3) ) WITH ( connector clickhouse, url jdbc:clickhouse://ck-host:8123, database-name profile, table-name user_profile, sink.batch-size 1000, sink.flush-interval 3s ); INSERT INTO clickhouse_user_profile SELECT user_id, member_level, register_time FROM mysql_user_profile;mysql-cdc连接器需要 MySQL 开启 binlog 且格式为 ROW。sink.batch-size设 1000 是攒批写入太小会频繁请求 ClickHouse太大延迟高。sink.flush-interval3 秒是兜底即使没攒够 1000 条也会刷。3.3 JDBC 连接器异常怎么排查热搜里「flink 的 jdbc 连接器异常」是高频问题常见三类现象原因解决No suitable driver found驱动 jar 没放进lib/或没在 SQL 里声明把对应 JDBC 驱动放到 Flinklib/目录重启集群Connection is not available连接池耗尽并发写入太高调大connection.max-retry-times或降低并行度Duplicate key写入失败维表有更新sink 没做主键去重ClickHouse 用ReplacingMergeTree引擎或 Flink 侧做last_value去重驱动问题最玄学很多时候是 Flink 集群和作业提交端 classpath 不一致本地跑通、集群报错血泪经验是统一把驱动放集群lib/。4. 避坑与排查实时画像链路的 5 个翻车现场4.1 现象画像标签延迟越来越高最后卡死原因状态后端用了默认的 HashMapStateBackend状态全在内存用户量一大就 OOM 或 GC 停顿。解决换成 RocksDBStateBackend并把状态 TTL 设上比如StateTtlConfig.newBuilder(Time.days(7))过期标签自动清理。4.2 现象窗口结果重复输出同一用户同一标签写了两遍原因用了处理时间窗口作业重启后重放数据窗口重新触发。解决改用事件时间窗口并开启EXACTLY_ONCEcheckpointsink 端做幂等HBase 用 RowKey 覆盖ClickHouse 用 ReplacingMergeTree。4.3 现象Kafka 消费积压但 Flink 并行度已经拉满原因keyBy 后数据倾斜某个热门用户的行为全进一个 subtask。解决keyBy 前加盐比如user_id _ random(0, 10)聚合后再去盐二次聚合。或者对超大 key 单独处理。4.4 现象SpringBoot 查画像接口 P99 超过 500ms原因每次请求都穿透到 HBase没加缓存。解决Redis 做一级缓存TTL 设 510 分钟HBase 做二级存储。热点用户画像可以本地 Caffeine 缓存但要注意多实例一致性。4.5 现象Flink 作业重启后画像数据错乱原因checkpoint 没配持久化路径或者 savepoint 没做。解决checkpoint 存 HDFS/S3升级作业用 savepoint 恢复别直接 kill 重启。state.checkpoints.dir一定要配。5. SpringBoot 整合 Flink 做画像查询服务的进阶技巧5.1 查询服务怎么和 Flink 作业解耦Flink 作业负责算SpringBoot 负责查两者通过存储层解耦不要用 RPC 直连。常见做法是 Flink 写 HBase/RedisSpringBoot 读。这样 Flink 作业重启不影响查询服务查询服务扩容也不影响计算。// SpringBoot 查画像Redis 优先HBase 兜底 public UserProfile getProfile(Long userId) { String key profile: userId; String cached redisTemplate.opsForValue().get(key); if (cached ! null) { return JSON.parseObject(cached, UserProfile.class); } // 缓存未命中查 HBase UserProfile profile hbaseDao.get(userId); if (profile ! null) { redisTemplate.opsForValue().set(key, JSON.toJSONString(profile), 10, TimeUnit.MINUTES); } return profile; }缓存 TTL 10 分钟是平衡实时性和 Redis 压力的常见值。如果画像更新频率高可以缩短到 12 分钟或者用 Flink 侧主动失效缓存。5.2 推荐引擎怎么消费画像推荐引擎拿画像做召回和排序的特征。召回阶段用画像标签圈人群比如「母婴高意向」用户走母婴商品池排序阶段把画像标签作为特征喂给模型。这里的关键是画像标签要版本化每次更新带一个version字段推荐引擎按版本取避免读到写了一半的画像。5.3 一个验证画像实时性的小技巧想验证画像是不是真的实时可以造一条测试埋点用测试账号点 5 次某类目商品然后立刻查画像接口看标签有没有在 1 分钟内更新。如果没有先查 Kafka 有没有积压再查 Flink 窗口有没有触发最后查存储写入有没有延迟。这个链路排查顺序我踩过很多次坑按这个顺序走基本能定位。我自己的习惯是任何实时链路上线前先跑一遍「埋点→Kafka→Flink→存储→接口」的端到端延迟打点把每一段的耗时记下来出问题时直接对比基线。这套画像推荐系统值不值得做取决于你的业务对时效的要求——如果 T1 画像已经够用上 Flink 就是过度设计如果用户行为变化快、推荐需要秒级响应那这条链路就是刚需。希望帮到你。本文还有配套的精品资源点击获取
返回列表