
简介一份基于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导致后面想算“用户当天浏览了多少商品”都算不出来。我一般会要求埋点至少要包含这几个字段字段示例作用userId1000234用户唯一标识itemIdSKU-88392商品IDbehaviorclick / cart / order / pay行为类型categoryId1203商品类目用于类目偏好timestamp1717300000000事件时间毫秒sessionId7f8a2c会话标识用于会话级统计deviceandroid / ios / pc渠道维度Kafka 的 Topic 设计不建议只用一个“user_behavior”大而全的 Topic。因为不同行为的吞吐量差异很大点击量可能是下单量的几百倍混在一起容易让下游的消费能力互相拖累。常见做法是拆成三个 Topicuser_click、user_cart、user_order。这样 Flink 可以针对不同 Topic 设置不同的并行度和 checkpoint 间隔。如果订单数据量小甚至可以一小时 check 一次点击数据量大就得每 30 秒 check 一次。生产环境里埋点数据还会经过一层 Nginx 日志采集或者由 SDK 直接推送到 Kafka。这里有一个关键参数acksall和retries3。很多团队为了吞吐把 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); DataStreamString rawStream env.addSource( new FlinkKafkaConsumer(user_click, new SimpleStringSchema(), kafkaProps) );这段代码里有几个参数值得注意。enableCheckpointing(30000, EXACTLY_ONCE)保证了从 Kafka 读取的数据不会因为故障而重复或丢失这是画像系统“准”的前提。auto.offset.resetearliest表示首次启动时从最早的 offset 开始读这样即使前一天链路挂掉第二天修复后也能把缺失的数据补回来。enable.auto.commitfalse配合 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 的流里。落表设计上我建议这们分。统计标签和规则标签直接算出来写入画像宽表算法标签往往是离线先跑再把结果同步到 RedisFlink 在实时流里读取 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 代码里要定义事件序列和超时时间。PatternClickEvent, ? pattern Pattern .ClickEventbegin(click) .where(event - event.getBehavior().equals(click)) .next(cart) .where(event - event.getBehavior().equals(cart)) .within(Time.minutes(5)); DataStreamClickEvent matched CEP.pattern(events, pattern) .select((MapString, 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 tableHBase 是列族存储加列不需要改 schema。HBase 的宽表设计RowKey 我推荐直接用倒序的 userId。比如 userId 是 1000234RowKey 就存4322001。因为 HBase 的 RowKey 是字典序排列正序的话相邻 userId 会落在同一个 Region热点问题严重倒序后数据分散到多个 Region。当然更规范的方案是加盐比如String.format(%02d_%s, userId % 100, userId)但加盐会牺牲顺序扫描能力。如果查询场景永远是“给定 userId 查全量标签”加盐也够用。表结构可以设计成三个列族列族包含标签类型示例列stats统计标签view_cnt_7d, order_amt_30drule规则标签is_vip, is_churn_riskalg算法标签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 RichSinkFunctionTuple2String, MapString, 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(Tuple2String, MapString, String tuple, Context context) throws Exception { Put put new Put(Bytes.toBytes(reverseUserId(tuple.f0))); for (Map.EntryString, 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.managedtrue让 Flink 统一管理 RocksDB 的内存同时限制每个 slot 的taskmanager.memory.managed.fraction0.4。另外别忽略state.backend.rocksdb.writebuffer.count写频繁时适当调低否则内存碎片严重。5.5 现象Flink CDC 同步业务库时数据不一致用户画像系统需要实时同步 MySQL 里的用户注册信息、订单状态。用 Flink CDC 没问题但一开始直接全程使用scan.incremental.snapshot.enabledtrue后发现某些 update 操作没有同步过来。原因是 CDC 底层读取 binlog如果 MySQL 的binlog_row_image是MINIMALupdate 事件里只包含被修改的列而我们的 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实现规则动态更新。MapStateDescriptorString, String ruleState new MapStateDescriptor( rule-config, BasicTypeInfo.STRING_TYPE_INFO, BasicTypeInfo.STRING_TYPE_INFO ); DataStreamString ruleStream env.addSource(new FlinkKafkaConsumer(rule_update, new SimpleStringSchema(), ruleProps)) .broadcast(ruleState); DataStreamTuple2String, String tagged clickStream .connect(ruleStream) .process(new BroadcastProcessFunctionClickEvent, String, Tuple2String, String() { Override public void processElement(ClickEvent value, ReadOnlyContext ctx, CollectorTuple2String, 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, CollectorTuple2String, 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 大小都要依据你的数据量和资源反复压测。希望上面这些从编码到排错的思路能帮到你让你在做这套画像系统时少熬夜、不翻车。本文还有配套的精品资源点击获取