简介:一份基于Flink流处理引擎的电商平台用户画像系统设计源码,面向大数据开发工程师与Java后端学习者,解决亿级电商数据实时处理与用户画像构建问题。压缩包共282个文件,含129个Java类、116个Java源文件,以及properties、XML、YAML等配置文件,其中dic字典文件与Kotlin模块文件辅助数据处理与扩展功能,整体约9.83MB。系统覆盖ViewService、InfoInService、RegisterCenter、PortraitAnalysis等模块,涵盖用户信息采集、注册存储、行为分析与画像生成全流程,并通过模块化设计提升可维护性。配套的README.md提供架构设计、接口定义与部署说明,便于快速上手。已有323人学习下载,适合希望掌握Flink实时计算、用户画像工程落地及大数据项目结构的开发者参考。
1. 为什么电商用户画像系统会选 Flink 而不是 Spark Streaming
做过实时画像的人都知道,画像系统的核心难题不是“算出来”,而是“算得及时”和“算得准”。用户刚刚点击了一个商品,下一秒推荐位就要反映出来;用户连续三次加购未支付,营销系统就要触发优惠券。这套逻辑如果靠离线批处理,跑完天都亮了,更别提应对大促时的流量洪峰。我经手过的电商画像项目,最初也试过 Spark Streaming,但真正上线后发现:精准的窗口计算、事件级的状态管理、以及和 Kafka、HBase 的生态衔接,Flink 明显更顺手。这个标题里的“Flink流处理引擎”和“用户画像系统”放在一起,本质上是在解决一个实时特征生产的问题——把用户每一次点击、搜索、下单、加购行为,在秒级延迟内变成标签,落到可查询的存储里。适合正在做实时数仓、推荐系统特征层、或者营销中台的开发者参考。下面我会从链路设计、代码骨架、存储选型、踩坑记录这四块,把一个可以照着改的源码级方案拆开讲清楚。
2. 用户画像系统的数据管道:从埋点到 Kafka 再到 Flink 的链路设计
2.1 埋点日志的字段设计与 Kafka Topic 划分
画像系统的最上游是埋点。很多团队在埋点阶段就偷懒,只采集了 userId、itemId、action,却漏掉了 timestamp、sessionId、deviceId,导致后面想算“用户当天浏览了多少商品”都算不出来。我一般会要求埋点至少要包含这几个字段:
| 字段 | 示例 | 作用 |
|---|---|---|
| userId | 1000234 | 用户唯一标识 |
| itemId | SKU-88392 | 商品ID |
| behavior | click / cart / order / pay | 行为类型 |
| categoryId | 1203 | 商品类目,用于类目偏好 |
| timestamp | 1717300000000 | 事件时间毫秒 |
| sessionId | 7f8a2c | 会话标识,用于会话级统计 |
| device | android / ios / pc | 渠道维度 |
Kafka 的 Topic 设计不建议只用一个“user_behavior”大而全的 Topic。因为不同行为的吞吐量差异很大,点击量可能是下单量的几百倍,混在一起容易让下游的消费能力互相拖累。常见做法是拆成三个 Topic:user_click、user_cart、user_order。这样 Flink 可以针对不同 Topic 设置不同的并行度和 checkpoint 间隔。如果订单数据量小,甚至可以一小时 check 一次;点击数据量大,就得每 30 秒 check 一次。
生产环境里,埋点数据还会经过一层 Nginx 日志采集,或者由 SDK 直接推送到 Kafka。这里有一个关键参数:acks=all和retries=3。很多团队为了吞吐把 acks 设为 1,结果 Kafka Broker 重启时日志丢失,画像标签就缺了一大块。既然做画像,数据完整性比那几百毫秒的延迟更重要。
2.2 Flink 消费 Kafka 的最小可运行代码与参数说明
用 Flink 消费 Kafka 是整套系统的基础。下面这个代码骨架是从我维护的画像项目里抽出来的最小版本,你可以直接抄下来改改就能跑。
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); env.enableCheckpointing(30000, CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(10000); env.getCheckpointConfig().setCheckpointTimeout(60000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); Properties kafkaProps = new Properties(); kafkaProps.setProperty("bootstrap.servers", "kafka-1:9092,kafka-2:9092"); kafkaProps.setProperty("group.id", "user-profile-group"); kafkaProps.setProperty("auto.offset.reset", "earliest"); kafkaProps.setProperty("enable.auto.commit", "false"); DataStream<String> rawStream = env.addSource( new FlinkKafkaConsumer<>("user_click", new SimpleStringSchema(), kafkaProps) );这段代码里有几个参数值得注意。enableCheckpointing(30000, EXACTLY_ONCE)保证了从 Kafka 读取的数据不会因为故障而重复或丢失,这是画像系统“准”的前提。auto.offset.reset=earliest表示首次启动时从最早的 offset 开始读,这样即使前一天链路挂掉,第二天修复后也能把缺失的数据补回来。enable.auto.commit=false配合 checkpoint,防止 offset 提交和数据处理不一致。
但这里有个大坑:如果你直接消费user_click原始 Topic,所有事件都会进来,包括爬虫和测试流量。更稳的做法是先在 Kafka 前面加一层“数据清洗”的 Topic——user_click_clean,由另一个 Flink 作业或者 Logstash 做过滤,画像作业只消费清洗后的数据。我见过有人把过滤逻辑写在画像主链路里,结果一个异常字段就导致整个作业重启。
另外,别把SimpleStringSchema直接用在生产。它只做字符串转换,不处理 JSON 解析。数据进 Flink 后应该立刻转成 POJO 或者 Avro。我在生产里用的是自定义的ClickEventDeserializationSchema,里面用 Jackson 解析,并且把字段缺失的情况兜底成默认值。
3. 画像标签计算:把原始行为加工成可查询的标签体系
3.1 标签模型:统计标签、规则标签、算法标签怎么落表
画像标签不是简单地把行为 count 一下就完事。它分三层:统计标签、规则标签、算法标签。统计标签是“用户最近 7 天点击次数”“最近 30 天下单金额”,直接基于窗口聚合;规则标签是“高价值用户”“流失预警用户”,规则写在代码里,比如“7 天未登录且曾经 30 天内下单超过 3 次”;算法标签是“性别预测”“购买力等级”,需要跑模型,但模型输出的结果也会灌进 Flink 的流里。
落表设计上,我建议这们分。统计标签和规则标签直接算出来写入画像宽表;算法标签往往是离线先跑,再把结果同步到 Redis,Flink 在实时流里读取 Redis 的模型结果,合并进标签。这样避免在 Flink 里跑复杂模型推理,把实时作业的稳定性保住了。
宽表的字段命名要统一。比如user_id、tag_name、tag_value、tag_time。不要用中文,不要在同一个表里有的字段叫cnt有的叫count。后面做特征查询时能少改很多代码。
3.2 使用 Flink SQL + CEP 计算实时标签的示例
Flink SQL 很适合做统计标签,因为声明式写法天然支持窗口聚合。下面这段 SQL 是“最近 1 小时用户点击类目 TOP3”的计算逻辑。
CREATE TABLE click_events ( user_id BIGINT, category_id BIGINT, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'user_click_clean', 'properties.bootstrap.servers' = 'kafka-1:9092', 'format' = 'json' ); CREATE TABLE category_preference ( user_id BIGINT, category_id BIGINT, click_cnt BIGINT, window_end TIMESTAMP(3), PRIMARY KEY (user_id, category_id) NOT ENFORCED ) WITH ( 'connector' = 'hbase-2.2', 'table-name' = 'user_profile:category_preference', 'zookeeper.quorum' = 'hbase-zk:2181' ); INSERT INTO category_preference SELECT user_id, category_id, COUNT(*) AS click_cnt, HOP_END(event_time, INTERVAL '5' MINUTE, INTERVAL '1' HOUR) FROM click_events GROUP BY user_id, category_id, HOP(event_time, INTERVAL '5' MINUTE, INTERVAL '1' HOUR);这段 SQL 的重点在于WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND。它允许事件乱序到达,最多等 5 秒。真实环境里用户手机网络无信号,事件可能延迟十几秒才上报,如果 watermark 设得太短,就会丢数据。但设置太长又会增加延迟,需要根据在线率调。
规则标签用 Flink CEP 写更合适。比如“用户点击了 A 商品后,5 分钟内加购了 B 商品,但最终没有下单”,这是典型的营销触发场景。CEP 代码里要定义事件序列和超时时间。
Pattern<ClickEvent, ?> pattern = Pattern .<ClickEvent>begin("click") .where(event -> event.getBehavior().equals("click")) .next("cart") .where(event -> event.getBehavior().equals("cart")) .within(Time.minutes(5)); DataStream<ClickEvent> matched = CEP.pattern(events, pattern) .select((Map<String, ClickEvent> map) -> { ClickEvent click = map.get("click"); ClickEvent cart = map.get("cart"); return new MarketingEvent(click.getUserId(), "click_then_cart", cart.getTimestamp()); });CEP 的.within(Time.minutes(5))定义了时间窗,超过 5 分钟这个序列就不算匹配。实际业务里,这个窗口值要根据品类决定。卖家电的决策周期长,可能是 7 天;卖零食的可能只有 20 分钟。不要把规则写死在代码里,建议把窗口时间放到配置中心或者 MySQL,用 BroadcastStream 动态更新。
到这里,标签已经算出来了,下一步就是怎么存储。这是画像系统最容易出问题的地方,下一章专门讲。
4. 画像存储与查询:为什么选 HBase + Redis 双层架构
4.1 标签宽表设计与 RowKey 设计要点
画像标签的存储,业界最常见的组合是 HBase + Redis。HBase 负责全量、按用户维度查询的宽表;Redis 负责实时性要求高的热数据,比如“当前用户是否是高价值用户”这种需要毫秒级返回的标签。为什么用 HBase 不用 MySQL?标签字段动辄几百个,MySQL 加一列就要 lock table,HBase 是列族存储,加列不需要改 schema。
HBase 的宽表设计,RowKey 我推荐直接用倒序的 userId。比如 userId 是 1000234,RowKey 就存4322001。因为 HBase 的 RowKey 是字典序排列,正序的话相邻 userId 会落在同一个 Region,热点问题严重;倒序后数据分散到多个 Region。当然更规范的方案是加盐,比如String.format("%02d_%s", userId % 100, userId),但加盐会牺牲顺序扫描能力。如果查询场景永远是“给定 userId 查全量标签”,加盐也够用。
表结构可以设计成三个列族:
| 列族 | 包含标签类型 | 示例列 |
|---|---|---|
| stats | 统计标签 | view_cnt_7d, order_amt_30d |
| rule | 规则标签 | is_vip, is_churn_risk |
| alg | 算法标签 | gender_pred, purchase_power |
每个列下面存标签值,tag_time通过 HBase 的 timestamp 维度记录。不要每个标签都建一张表,查询时跨表 join 在 HBase 里是很痛苦的事。
4.2 Flink 写入 HBase 的 sink 代码与参数调优
Flink 写 HBase 最常见的姿势是继承RichSinkFunction,自己管理 BufferedMutator。直接调用 HBase 的 put 会有性能问题,因为每个 put 都是一次 RPC。
public class ProfileSink extends RichSinkFunction<Tuple2<String, Map<String, String>>> { private Connection conn; private BufferedMutator mutator; @Override public void open(Configuration parameters) throws Exception { org.apache.hadoop.conf.Configuration hbaseConfig = HBaseConfiguration.create(); hbaseConfig.set("hbase.zookeeper.quorum", "hbase-zk:2181"); conn = ConnectionFactory.createConnection(hbaseConfig); BufferedMutatorParams params = new BufferedMutatorParams(TableName.valueOf("user_profile:profile")); params.writeBufferSize(8 * 1024 * 1024); // 8MB buffer mutator = conn.getBufferedMutator(params); } @Override public void invoke(Tuple2<String, Map<String, String>> tuple, Context context) throws Exception { Put put = new Put(Bytes.toBytes(reverseUserId(tuple.f0))); for (Map.Entry<String, String> entry : tuple.f1.entrySet()) { put.addColumn(Bytes.toBytes("stats"), Bytes.toBytes(entry.getKey()), Bytes.toBytes(entry.getValue())); } mutator.mutate(put); if (mutator.getWriteBufferSize() > 16 * 1024 * 1024) { mutator.flush(); } } @Override public void close() throws Exception { if (mutator != null) mutator.close(); if (conn != null) conn.close(); } }这块有两个参数很关键。writeBufferSize(8 * 1024 * 1024)是缓冲区的阈值,到 8MB 就自动刷写。调太小会频繁 RPC,调太大会让内存压力上升,集群规模不大的话 8MB 比较稳。还有一个隐性参数是 HBase 服务端的hbase.client.write.buffer,客户端的 buffer 要和服务端匹配,否则会出现服务端迫不及待刷写,导致写放大。
另外,invoke里判断writeBufferSize > 16MB才 flush,这是为了手动控制刷写节奏,防止单条数据过大触发自动刷写时阻塞 Flink 主线程。因为mutator.mutate是异步的,flush 是同步的,频繁同步 flush 会拖慢吞吐。
写入 HBase 之前还有一个必要步骤:去重。Flink checkpoint 开启 EXACTLY_ONCE 后,Kafka 源不会重发,但 HBase sink 不支持事务性写入,如果上游手动重放数据,就可能导致重复。简单做法是在 HBase 表设计时把同一 userId 的标签列用putToSameCell,利用 HBase 的覆盖写特性,重复写入同一个列会覆盖旧值,天然幂等。这一点不用太担心。
5. Flink 用户画像系统避坑指南:5 个真实踩坑记录
5.1 现象:Kafka 消费延迟越来越高,但 CPU 没跑满
有一次我负责的画像作业,Kafka 堆积量从 100 万涨到 5000 万,Flink 监控面板显示 CPU 使用率只有 20%,但消费速率就是上不去。看线程 dump 发现大量线程阻塞在HBase.put上。原因是 HBase RegionServer 的 MemStore 达到阈值后触发了 flush,而客户端没有开启异步批量,每次 put 都等 RPC 返回。解决方法是改用 BufferedMutator,并且把 Flink 算子的并行度和 HBase Region 数量对齐,避免写倾斜。另外,检查是不是所有字段都写进了同一个列族,导致单 Region 写入压力过大。
5.2 现象:窗口计算出来的标签总是偏少
Flink 的滚动窗口统计“过去 1 小时点击量”,结果比业务方从数据库查出来的少了 30%。排查发现埋点日志里的timestamp是客户端时间,而不是服务器接收时间。用户手机时钟不准,或者客户端把事件缓存了几分钟,导致某些事件的时间戳晚于 watermark,直接被判定为迟到数据丢弃。解决方法是统一用 Kafka 的 ingestion time,或者在埋点 SDK 里强制在服务端接收时重打时间戳。我后来直接把EventTime改成了ProcessingTime,配合 5 秒的乱序容忍,才把数据补齐。
5.3 现象:HBase 写入 hotspot,部分 Region 数据量涨到其他 Region 的十倍
因为 RowKey 用的是 userId 正序,前 1000 个用户都是老用户,频繁下单,全部落在前几个 Region。大量写入都压在那几个 RegionServer 上,集群整体 CPU 不高,但部分节点告警。改成了倒序 RowKey 之后,数据分布立刻均匀了。如果倒序还不够,建议对 userId 做哈希取模加盐,比如userId % 200作为前缀,保证 200 个桶的分散度。
5.4 现象:状态后端 RocksDB 导致 OOM
画像作业用了 CEP 和大量窗口聚合,状态越积越大。默认的 HashMapStateBackend 放不下,换成了 RocksDBStateBackend,结果 JobManager 直接 OOM。原因是 RocksDB 的 block cache 和 write buffer 默认配置偏大,多个 slot 共享内存时叠加超限。解决方法是设置state.backend.rocksdb.memory.managed=true,让 Flink 统一管理 RocksDB 的内存,同时限制每个 slot 的taskmanager.memory.managed.fraction=0.4。另外别忽略state.backend.rocksdb.writebuffer.count,写频繁时适当调低,否则内存碎片严重。
5.5 现象:Flink CDC 同步业务库时数据不一致
用户画像系统需要实时同步 MySQL 里的用户注册信息、订单状态。用 Flink CDC 没问题,但一开始直接全程使用scan.incremental.snapshot.enabled=true后,发现某些 update 操作没有同步过来。原因是 CDC 底层读取 binlog,如果 MySQL 的binlog_row_image是MINIMAL,update 事件里只包含被修改的列,而我们的 JSON 解析器要求所有字段齐全。解决方法是把 MySQL 的binlog_row_image设为FULL,同时在 CDC 配置里加上debezium.event.deserialization.failure.handling为warn,先别让作业挂掉,通过日志排查问题。
这些坑背后都有一个共性:实时链路里每个环节都可能因为数据偏差导致最终标签不准。所以系统上线前一定要做数据质量验证,这也是最后一章要说的内容。
6. 把画像数据回灌到业务系统的进阶技巧:实时特征服务与验证方法
6.1 用 BroadcastStream 加载画像规则
前面提到的规则标签,如果每次改规则都要重启作业,就太被动了。Flink 的 BroadcastStream 可以解决这个问题。把规则配置放在一个 Kafka Topic 里,比如rule_update,然后通过broadcast和主数据流 connect,实现规则动态更新。
MapStateDescriptor<String, String> ruleState = new MapStateDescriptor<>( "rule-config", BasicTypeInfo.STRING_TYPE_INFO, BasicTypeInfo.STRING_TYPE_INFO ); DataStream<String> ruleStream = env.addSource(new FlinkKafkaConsumer<>("rule_update", new SimpleStringSchema(), ruleProps)) .broadcast(ruleState); DataStream<Tuple2<String, String>> tagged = clickStream .connect(ruleStream) .process(new BroadcastProcessFunction<ClickEvent, String, Tuple2<String, String>>() { @Override public void processElement(ClickEvent value, ReadOnlyContext ctx, Collector<Tuple2<String, String>> out) { String rule = ctx.getBroadcastState(ruleState).get("cart_in_5min"); if (rule != null && rule.equals("true")) { out.collect(new Tuple2<>(value.getUserId(), "is_interest")); } } @Override public void processBroadcastElement(String value, Context ctx, Collector<Tuple2<String, String>> out) { ctx.getBroadcastState(ruleState).put("cart_in_5min", value); } });注意BroadcastProcessFunction里processElement只能读广播状态,不能修改,修改只能发生在processBroadcastElement里。这是并发安全的硬性约束。很多同学想当然地在每条数据里写广播状态,运行时会直接抛异常。
6.2 画像质量校验:用 Redis 对比抽样验证
画像系统上线后,最难回答的问题是“你算的标签到底准不准”。我的习惯做法是用 Redis 做两组数据的对比。一组是 Flink 实时写入的标签,另一组是离线数仓每天凌晨算好的标签。对同一批 userId,从两处读出标签值,计算不一致率。
例如,离线统计“用户 30 天订单金额”是 2500 元,实时画像算出来是 2498 元,差 2 元可能是因为有一笔刚发生的订单还没进 Kafka,或者在窗口边界被切走了。不一致率控制在 2% 以内算正常,超过 5% 就要查链路。具体验证代码可以简单写一个定时任务:
// 伪代码:每天凌晨 1 点,对比 Redis 中实时标签和 Hive 中离线标签 for (String userId : sampleUserIds) { double realTimeValue = getFromRedis("profile:" + userId + ":order_amt_30d"); double offlineValue = getFromHive("select order_amt_30d from offline_profile where user_id = ?"); double diff = Math.abs(realTimeValue - offlineValue) / offlineValue; if (diff > 0.05) { logger.warn("user {} diff too large: realTime={}, offline={}", userId, realTimeValue, offlineValue); } }这个对比脚本不复杂,但能帮你在业务方投诉之前发现问题。我还习惯在 Flink 作业里加一个late_element_count计数器,每天看迟到数据量占总量的比例,超过阈值就调大 watermark 容忍时间。毕竟画像系统是给业务决策用的,一个坏的标签比没有标签更可怕。
最后提醒一句:Flink 作业的参数没有银弹,并行度、checkpoint 间隔、buffer 大小都要依据你的数据量和资源反复压测。希望上面这些从编码到排错的思路能帮到你,让你在做这套画像系统时少熬夜、不翻车。
本文还有配套的精品资源,点击获取