Flink源码实战:电商用户画像系统从行为日志到实时标签的工程剖析
简介基于Flink流处理引擎的电商用户画像系统源码面向大数据开发、Java后端及流计算学习者解决海量电商数据实时分析与用户画像构建难题。系统针对亿级数据场景通过Flink高效处理用户行为流提取关键特征为精准推荐、广告投放及用户粘性提升提供支撑。压缩包共282个文件、约9.83MB包含116个Java源文件、129个class文件、15个properties配置、9个XML及2个YAML配置、6个字典文件、4个Kotlin模块和1个说明文档目录结构清晰。源码覆盖数据收集、注册中心、画像分析等核心模块如InfoInService、RegisterCenter、PortraitAnalysis模块化设计便于二次开发通过阅读代码可掌握从数据接入、清洗、特征提取到画像标签生成的完整链路同时理解Flink窗口计算与状态管理等关键机制并参考README快速部署。目前已有323人学习适合作为Flink流处理项目实战参考尤其适合需要深入理解Flink实际落地的开发者代码结构清晰可直接改造用于电商场景。1. 从亿级行为日志到用户标签Flink 用户画像源码的工程切面用户画像系统听起来像“采集日志、打个标签”这么简单真去拆源码才发现难点全在工程细节上。这套基于Flink流处理引擎的电商用户画像系统源码总共282个文件、129个Java类把标签计算拆成了UserTypeTask、BrandLikeTask、ChaomanandwomenTask等一系列独立任务再配合KMeansRunbyusergroup做用户分群、Logistic做分类预测、MongodataControl统一读写MongoDB。也就是说画像不是一张表而是一组可独立调度的流作业。对正在做实时标签平台或者想把用户画像从离线批任务迁到实时流的团队这套源码里的任务划分和存储设计值得完整过一遍。2. 流计算选型与画像模块拆分Flink 的状态与窗口如何支撑标签任务2.1 为什么生产画像任务优先选 Flink 而不是 Spark Streaming用户画像的本质是把用户行为序列压缩成一组可检索的特征。电商场景下这个压缩过程要同时满足两个条件延迟在秒级到分钟级结果能按用户维度持续更新。Flink 在这类需求里的核心优势有三个。第一是状态管理机制KeyedStream 上的 ValueState 和 MapState 可以直接保存“用户当前累计的浏览时长、最近30天购买类目”这类中间结果不需要外部存储兜底。第二是事件时间与水印机制用户行为日志在链路里必然存在乱序基于 ingestion time 处理会把数据算错靠 event time 配合 watermark 才能得到准确划分的窗口结果。第三是精确一次语义checkpoint 配合 Kafka 的 offset 管理保证标签任务在重启后不重不漏。能力维度Spark Streaming微批Flink真流对画像场景的影响延迟秒级~分钟级由批间隔决定毫秒~秒级事件驱动实时标签更新频率状态管理依赖外部存储做增量内置状态后端RocksDB等用户累计特征是否易算乱序处理窗口内靠延迟数据修正watermark allowedLateness点击/下单日志乱序是否准确精确一次较晚版本支持原生支持广告计费类标签是否可信Flink 的状态同时是这套源码里很多 Task 能独立运行的前提。比如 BaijiaTask 这种按数据源拆出的任务入口在 Flink 里就是一个独立 Job 或独立算子链共享同一个集群资源但状态互不干扰。换成批处理引擎就得自己在外部存储里维护“这个任务算到哪了”复杂度完全不一样。2.2 从类名反推画像系统的模块边界拿到源码先别急着跑把 class 按职责归类基本就能还原出一张架构图。我按生产画像系统的分层习惯做了个映射源码类名推断职责对应分层ViewService、InfoInService、RegisterCenter用户信息录入、注册存储、前端展示服务接入与存储层PortraitAnalysis用户行为分析并构建画像的计算入口计算层入口UserTypeTask、ChaomanandwomenTask用户类型、潮男潮女等规则标签规则标签层BrandLikeTask品牌偏好统计偏好标签层KMeansRunbyusergroup、Logistic用户分群与分类预测算法模型层UserGroupInfo、UserGroupSecondMap分组 POJO 与二次映射结果表达层MongodataControlMongoDB 连接与读写控制存储层这套划分有一个明显特点标签计算任务和算法任务是解耦的。UserTypeTask 这类规则任务不依赖模型结果KMeans 分群输出又反过来作为 UserGroupSecondMap 的输入。同一个行为数据流可以被多个作业并行消费每个作业产出一类标签最后由 UserGroupInfo 汇总落库。这种“多作业、单汇聚”的结构比单作业里串行算几十个标签更好扩展也更好排查哪个标签没更新直接看对应 Job 的监控即可。2.3 画像主链路行为日志到 MongoDB 的实时管道用户画像的数据链路可以抽象成行为日志 → Kafka → Flink 标签作业群 → MongoDB。MongodataControl 的存在说明这套系统没有把画像结果直接推进 Redis而是统一落到 MongoDB再由下游服务读取。MongoDB 在这类场景里的好处是文档模型贴合“一个用户一条画像记录”的形态嵌套的标签数组和偏好对象可以原样存取不需要像关系库那样拆表。代价是更新频繁时写放大明显生产里一般会搭配 TTL 索引清理过期标签或者把热标签放 Redis、全量画像放 Mongo。提示如果本地只有源码没有数据源可以用 Flink 自带的 SocketTextStreamFunction 或者 Kafka connector 的测试模式先把链路跑通再去接真实业务日志。3. 标签计算与用户分群KMeans 聚类、规则标签和品牌偏好的 Flink 算子实现3.1 KMeans 分群离线训练模型 流式打标签KMeansRunbyusergroup 这个类名里有两个关键信息KMeans 算法和 by usergroup。生产里很少有人真的在流上跑迭代式 KMeans更常见的做法是用离线历史行为数据训练出 K 个聚类中心把中心点以广播状态BroadcastState下发到 Flink 作业每条用户特征进来后计算到各中心点的距离取最近的中心作为该用户的分群 ID。这个方案把训练和预测分离既保证聚类结果可离线评估又保证线上打标签的实时性。// 伪代码基于广播聚类中心做实时用户分群 public class KMeansRunbyusergroup extends KeyedProcessFunctionString, UserFeature, UserGroupInfo { // 广播状态保存离线训练好的 K 个聚类中心 private transient BroadcastStateString, double[] centers; Override public void processElement(UserFeature feature, Context ctx, CollectorUserGroupInfo out) throws Exception { double[] current centers.get(centers); if (current null) { // 聚类中心还没广播到位先走默认分组 out.collect(new UserGroupInfo(feature.getUserId(), default)); return; } int groupId nearestCenter(feature.toArray(), current); out.collect(new UserGroupInfo(feature.getUserId(), group_ groupId)); } private int nearestCenter(double[] point, double[] center) { // 计算欧氏距离返回最近中心的下标 double minDist Double.MAX_VALUE; int index -1; for (int i 0; i center.length; i) { double dist distance(point, center[i]); if (dist minDist) { minDist dist; index i; } } return index; } }这段伪代码里BroadcastState 是 Flink 专门为“配置/模型动态更新”设计的机制。广播状态里的数据对所有并行子任务可见而且可以在作业不重启的情况下更新。也就是说运营重新跑了一次离线聚类只需要把新中心点广播进去线上分群结果在下一轮特征进来时自然切换。如果有 5000 万用户特征要打分并行度调到 16 或 32吞吐量基本能压到单作业每秒几十万条瓶颈通常在 MongoDB 写入端而不是计算端。同一目录里的 Logistic 类则是训练好的 LR 模型用于购买意向这类二分类标签预测路径和 KMeans 完全一样模型参数走广播下发。聚类训练阶段有几个参数直接决定标签质量参数建议值说明K 值520按业务人群颗粒度选常用肘部法确定特征维度1050浏览/加购/下单频次标准化后入模训练数据窗口90 天覆盖足够长的行为周期中心点更新频率每天或每周广播更新不重启作业3.2 规则标签UserTypeTask 与 ChaomanandwomenTask 的处理骨架规则标签是画像系统里数量最多的一类。UserTypeTask 判断新客/老客/高活跃/流失风险ChaomanandwomenTask 判断用户是否属于潮男潮女人群逻辑上都是“读取行为特征 → 命中规则 → 输出标签”差别只在规则条件和特征来源。这类任务在 Flink 里最稳的写法是 KeyedProcessFunction 配合定时器按用户维度累积事件而不是每个事件单独打标签否则同一个用户在一个小时内会被打上互相矛盾的标签。// 基于 30 天窗口的用户类型判定 public class UserTypeTask extends KeyedProcessFunctionString, UserLog, TagOutput { private transient ValueStateLong firstSeenTs; private transient ValueStateLong lastSeenTs; Override public void processElement(UserLog log, Context ctx, CollectorTagOutput out) throws Exception { if (firstSeenTs.value() null) { firstSeenTs.update(log.getTs()); // 注册 30 天后的定时器到期输出一次判断结果 ctx.timerService().registerEventTimeTimer(log.getTs() 30L * 24 * 3600 * 1000); } lastSeenTs.update(log.getTs()); } Override public void onTimer(long timestamp, OnTimerContext ctx, CollectorTagOutput out) throws Exception { long first firstSeenTs.value(); long last lastSeenTs.value(); String type (last - first 30L * 24 * 3600 * 1000) ? loyal : new; out.collect(new TagOutput(ctx.getCurrentKey(), type)); } }定时器在这里的作用是延迟计算。有一类标签必须等时间窗口结束才能下结论比如 30 天未访问才算流失。用事件时间定时器Flink 会在 watermark 越过指定时间点时触发计算天然规避了“每条日志都重新判断一次”的浪费。firstSeenTs 和 lastSeenTs 是 KeyedStream 上的 ValueState同一个用户的所有事件都会落到同一个 key 上状态互相可见。需要调优的参数有两个state.backend 选 RocksDB 还是内存后端以及 state.ttl。这类中间状态如果不设 TTL会随着用户量增长把磁盘撑爆生产上我一般给规则标签中间态设 35 天 TTL刚好比 30 天窗口多 5 天冗余。3.3 BrandLikeTask用窗口聚合统计品牌偏好品牌偏好是电商画像里最典型的可聚合标签。用户在 7 天内看了多少次某品牌商品、加购几次、下单几次换算成品牌偏好分。实现上最直接的是事件时间窗口加聚合// 7 天品牌偏好统计输出 (userId, brandId, score) DataStreamUserLog logs ...; logs.keyBy(log - log.getUserId() _ log.getBrandId()) .window(TumblingEventTimeWindows.of(Time.days(7))) .aggregate(new AggregateFunctionUserLog, BrandAccumulator, BrandScore() { Override public BrandAccumulator createAccumulator() { return new BrandAccumulator(); } Override public BrandAccumulator add(UserLog log, BrandAccumulator acc) { // 浏览计 1 分加购计 3 分下单计 5 分 acc.addScore(log.getActionType().score()); return acc; } Override public BrandScore getResult(BrandAccumulator acc) { return new BrandScore(acc.getUserId(), acc.getBrandId(), acc.getScore()); } Override public BrandAccumulator merge(BrandAccumulator a, BrandAccumulator b) { a.merge(b); return a; } });窗口长度、分数权重、keyBy 粒度是这三个关键调优点。keyBy 粒度建议用 userId brandId 联合 key而不是只按 userId 分后者会把同一个用户的所有品牌偏好塞进同一个算子子任务数据倾斜时个别子任务会成为瓶颈。分数权重可以直接做成配置从 properties 文件里读运营调权重不需要改代码重新打包。窗口聚合结果默认逐条输出如果下游 MongoDB 写入压力大可以在聚合算子后面加一个 TumblingProcessingTimeWindows(Time.seconds(5)) 做微批攒批把写入吞吐降下来。这种“计算用事件时间、输出用处理时间攒批”的写法是 Flink 画像任务上比较实用的性能优化手段。3.4 UserGroupInfo 与 UserGroupSecondMap标签结果的组合作业最后一个环节是把分散的标签合并到用户维度。UserGroupInfo 是携带 userId 和 groupId 的 POJOUserGroupSecondMap 字面上看是做第二次映射。常见做法是从 MongoDB 读取用户的已有画像把本次新算出的标签合并进文档再写回。合并过程用 CoGroup 或 Connect 操作把 UserTypeTask、BrandLikeTask、KMeans 三个作业的输出按 userId 连接起来。需要注意多流 join 在流处理里是有代价的状态会随用户量线性增长。如果标签种类很多更实际的做法是每条标签写一个独立字段到 MongoDB用 $set 按字段更新而不是把全部标签拼成一个大对象再覆盖写。4. 配置体系与任务提交从 properties 到 Flink SQL 和 MongoDB 连接4.1 三类配置文件的分工这套源码里有 15 个 properties、9 个 XML、2 个 YAML乍一看有点杂其实是典型的多层配置结构。YAML 管服务级配置properties 管环境相关的连接信息XML 管日志和 mapper 映射。源码里还带了 4 个 Kotlin 模块文件大概率是 Gradle Kotlin DSL 的构建脚本不影响运行时不用管它。我习惯按“环境无关 / 环境相关”分开管理文件类型典型内容变更频率注意点YAML服务端口、数据源、Flink 运行模式低多环境用 profile 拆分propertiesKafka 地址、MongoDB URI、Redis 密码中禁止提交到 Git用配置中心XMLlog4j2、mybatis mapper低注意 classpath 冲突把连接信息放 properties 而不是写死在 Java 类里最大的好处是换环境不用改代码。MongodataControl 读取 MongoDB 连接串时如果源码里没用配置中心常见的做法是读取 classpath 下的 mongodb.properties里面配 mongodb.urimongodb://user:passhost:27017/portrait。生产环境里我会再套一层环境变量覆盖比如 MONGO_URI 存在的话优先用环境变量否则读本地文件这样 Docker 部署时不用重新构建镜像。4.2 Flink 任务提交与资源参数本地开发直接在 IDE 里跑 main 方法生产则是打包后提交到集群。下面是一个标准提交命令假设任务主类是 com.portrait.BrandLikeTask。flink run -m yarn-cluster \ -ys 2 \ -yjm 1024m \ -ytm 2048m \ -p 16 \ -d \ -c com.portrait.BrandLikeTask \ portrait-flink-1.0.jar参数说明-ys 指定每个 TaskManager 的 slot 数决定并行度上限-yjm 和 -ytm 分别是 JobManager 和 TaskManager 内存-p 是作业并行度要与上游 Kafka 分区数对齐。比如 Kafka 分区 12并行度设 12 或 24 都没有问题设 48 就可能让部分子任务空转-d 表示 detached 后台运行。如果任务跑一段时间后出现背压优先调并行度和窗口聚合的攒批逻辑而不是无脑加内存画像任务大多是 IO 型内存堆再大MongoDB 写不进去照样背压。提示任务提交后很快失败先看日志。画像任务启动失败常见原因是 MongoDB driver 和 Flink 内置的 Netty 版本冲突用 mvn dependency:tree 检查依赖树最直接。4.3 Flink SQL 与 JDBC 连接器标签查询与维表关联画像结果落库后往往还要做标签查询Flink SQL 是绕不开的一环。用 Flink SQL 建一张映射 MongoDB 的维表可以在流作业里直接通过用户 ID 关联画像标签避免手写 MongoDB 查询CREATE TABLE user_profile ( user_id STRING, brand_tags ARRAYSTRING, group_id STRING, update_time TIMESTAMP(3), PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( connector mongodb, uri mongodb://localhost:27017/portrait, database portrait, collection user_profile, lookup.cache.max-rows 50000, lookup.cache.ttl 1 h );有了这张维表实时订单流就可以在 SQL 里直接 JOIN 画像标签做人群过滤。lookup.cache 两个参数是关键max-rows 控制缓存条目上限ttl 控制缓存过期时间。画像标签这种低频变化的维表缓存 TTL 设 1 小时能省掉绝大多数 MongoDB 查询。用 sql-client 或者 SQL Gateway 提交这份 DDL和 DataStream API 作业共享集群资源。如果你在这个 JOIN 上报错最常见的是 JDBC 连接器异常MongoDB connector 用的是异步 IO连接数不够时会出现 MongoSocketReadTimeout。处理思路是调大连接池参数把 maxConnections 从默认值往上提同时把 socketTimeout 设到 5 秒以上。4.4 水印与乱序Flink SQL 里如何声明事件时间电商日志经过多个网关到 Kafka乱序是常态。订单日志晚到 10 秒是正常波动晚到 5 分钟也时有发生。Flink SQL 里用 WATERMARK 语法声明事件时间CREATE TABLE user_behavior ( user_id STRING, brand_id STRING, action_type STRING, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 2 MINUTE ) WITH ( connector kafka, topic user_log, properties.bootstrap.servers kafka:9092, properties.group.id portrait-group, format json );WATERMARK FOR ts AS ts - INTERVAL 2 MINUTE 表示允许 2 分钟内的乱序数据。这个值设得太小晚到数据会被丢弃设得太大窗口结果输出会延迟。对画像任务来说标签晚两分钟出并不致命丢数据才是致命的丢掉的点击行为会导致偏好分数永远偏低。实际项目里我会先用 Metrics 观察日志时间戳和当前时间的差值分布取 95 分位作为 watermark 延迟预算再留 20% 余量。flink 菜鸟教程里那套“固定延迟加观察调整”的方法放在这里要强调一点不要照抄别人配的 10 秒你的日志链路时延分布跟别人不一样。4.5 MongodataControl 的读写封装最后一个关键类是 MongodataControl。它的职责是把 MongoDB 的增删改查统一封装避免每个 Task 都直接持有 MongoClient。常见实现是单例 MongoClient 加带超时时间的写操作。注意一个细节MongoClient 是线程安全的但 MongoCollection 写并发很高时建议用 bulkWrite 攒批写入而不是逐条 insert。画像标签这类小文档单条 insert 的网络往返开销甚至比写入本身还大攒批后吞吐能提升一个数量级。// 攒批写入画像标签到 MongoDB public class MongodataControl { private static final int BATCH_SIZE 200; private final ListWriteModelDocument buffer new ArrayList(); public synchronized void upsert(String userId, String field, Object value) { Document filter new Document(user_id, userId); Document update new Document($set, new Document(field, value)); buffer.add(new UpdateOneModel(filter, update, new UpdateOptions().upsert(true))); if (buffer.size() BATCH_SIZE) { flush(); } } public synchronized void flush() { if (buffer.isEmpty()) return; collection.bulkWrite(buffer, new BulkWriteOptions().ordered(false)); buffer.clear(); } }这里选 UpdateOneModel 而不是 ReplaceOneModel对应之前说的字段级更新。ordered(false) 让写入不按顺序串行执行多个 Task 并发写不同字段时互不阻塞。BATCH_SIZE 一般设 100 到 500太小攒批效果不明显太大单次请求内存占用高且失败重放代价大。如果写 MongoDB 的算子出现背压不要先加并行度先看 bulkWrite 批次有没有触达 200再看 MongoDB 服务端的连接数与 CPU下游先看上游再看。5. 从实时标签到实时数仓Flink CDC、数据血缘与画像结果交付5.1 画像标签的两种消费方式标签算出来后消费方式决定存储选型。推荐和广告系统要求毫秒级读取MongoDB 扛不住高 QPS 的按用户查询这时候把热标签同步一份到 Redis 是常见做法BI 分析要看人群整体分布则需要把画像结果导出到 ClickHouse 或 Hive 做 OLAP。Flink 可以同时做这两件事一个作业写 MongoDB另一个作业读 MongoDB 的变更流写 Redis。从这个角度看Flink 的使用场景不只在计算层数据分发层同样能承担。5.2 引入 Flink CDC 更新画像源数据画像准确度取决于源数据质量。过去只能等 T1 批任务修正标签现在用 Flink CDC 监听业务库变更可以近实时更新用户的基础属性。社区里经常讨论 Flink 2.2.1 搭配 Flink CDC 3.5.0 的 Docker 部署组合核心思路是先用 Debezium 捕获 binlog再通过 Flink CDC connector 把变更数据以 Changelog 流方式接入画像作业直接修正 UserGroupInfo 里的基础字段不需要重算全量标签。CDC 接进来的数据对 watermark 要特别注意binlog 时间戳和业务时间往往有偏差不能直接拿它算窗口。5.3 数据血缘让标签可追溯画像标签的可解释性越来越重要用户问为什么给他推这个品牌要能回答“因为过去 7 天浏览了 12 次该品牌”。Flink 作业如果都是 SQL 写的可以直接利用 Flink 的 lineage 能力从 Catalog 拿表级血缘纯 DataStream API 的任务则需要自己埋点输出血缘信息。最简单的做法是每个 Task 输出时附带一个 sourceVersion 字段记录这次标签计算用的 Kafka topic 和时间范围MongoDB 里每条标签旁边存这个版本号追溯时查一下就能定位到数据来源。5.4 上线前必查的几个验证点上线前我会逐项确认并行度和上游分区数对齐状态 TTL 已经配置MongoDB 连接池上限和 bulkWrite 批次合理watermark 延迟基于真实日志分布设定CDC 任务单独分配资源避免和核心标签作业抢 slot。验证方法可以简单粗暴拿最近 7 天真实日志回放对比新实时任务产出的标签和离线批任务产出的标签一致率低于 95% 就先查 watermark 和状态清理逻辑再查维表缓存是否命中过期数据。这套回放验证法比看几个样例准得多。本文还有配套的精品资源点击获取