
简介面向大数据开发与实时计算场景这套基于Flink流处理引擎的电商用户画像系统源码专为处理亿级用户行为数据、构建精准画像而设计。系统以Java为核心覆盖ViewService、InfoInService、RegisterCenter、PortraitAnalysis等模块提供从数据接入、信息注册到画像分析的全链路实现适合理解Flink在电商推荐、精准营销与用户增长中的落地方式。包体共282个文件、约9.83MB包含116个Java源文件与129个class文件另有properties、xml、yml配置项、dic字典及kotlin_module模块配置与字典便于灵活调整运行参数Java源码则可直接用于二次开发和学习调用链。项目附带README文档说明架构设计、依赖关系与部署步骤降低上手门槛。目前已有323人学习浏览。对希望掌握用户画像系统搭建、Flink流处理项目实践或电商平台个性化推荐实现细节的开发者均可从中提取模块拆分思路、配置管理方式和实时计算代码骨架用于毕业设计、课程项目或生产原型参考参考价值直接。1. 为什么用户画像系统要上流处理引擎先说一个反直觉的结论离线画像跑得再准对大促当天的运营来说依然是马后炮。用户在 10:25 把商品加进购物车10:32 犹豫退出10:35 领了一张券又回来——这一连串行为发生在分钟甚至秒级尺度内离线 T1 画像根本来不及感知。电商做用户画像系统最终目的是为了让推荐、营销、客服在一个「刚刚发生」的时间窗口内用上标签而不是昨天算好的旧档案。所以「基于 Flink 流处理引擎的电商平台用户画像系统」要解决的核心问题只有两个一是把散落在订单、点击、浏览、收藏日志里的用户行为实时聚合成标签二是让这些标签可以被后续服务以低延迟方式查询。适合读这篇文章的人是已经在做离线画像、准备往实时架构演进或者刚被业务方要求「大促前一天晚上给我出实时人群包」的工程师。接下来的内容会围绕标签体系设计、Flink 实现路径、存储选型和端到端验证展开代码部分可以直接落地不掺假。2. 先想清楚画像模型再动 Flink 代码2.1 标签分类直接决定计算引擎的形态实时用户画像系统里标签通常按计算逻辑分成三类统计型、规则型、算法型。统计型标签最容易理解比如「近 7 天加购次数」「近 1 小时浏览商品数」本质是对行为流做窗口聚合Flink 的滚动窗口配合状态后端可以秒级产出。规则型标签要麻烦一点比如「高活跃但近 3 天未下单」这类标签需要同时看到长期状态和短期行为意味着 Flink 作业必须维护跨窗口的用户状态。算法型标签则是对行为序列做打分或分类典型场景是「购买意向分」「流失概率分」。这类逻辑如果完全放在 Flink 里跑模型推理会比较重常见做法是让 Flink 将实时特征拼接好后调用外部推理服务再把打分结果写回标签存储。前期不建议在 Flink 内部做复杂机器学习一是调试困难二是 Flink 的状态后端处理高维特征容易导致检查点膨胀。设计阶段就要把这三种标签分开建模因为它们在存储、时效和召回策略上完全不同。建议用一张标签元数据表对运营和技术人员统一暴露字段定义表结构按「标签编码、标签名称、标签类型、值类型、更新频率、所属业务线」组织。值类型决定了下游怎么用——布尔型适合圈选人群数值型适合排序枚举型适合内容匹配。2.2 标签存储选型一张表装不下所有画像画像数据的访问模式有两个明显特征写入是流式的单次写很小但频率高读取是点查或者范围扫要求毫秒级响应。单一存储没法同时应付这两个特征常见做法是分层存储组件存储内容时效性查询场景Redis秒级更新的高频标签如实时活跃状态秒级在线推荐接口读取HBase全量用户标签宽表Rowkey 为用户 ID分钟级人群圈选、离线同步Elasticsearch可检索的标签索引分钟级运营条件组合筛选人群MySQL / TiDB标签元数据、维度配置小时级标签管理后台Redis 里存的是「当前这一刻」的状态例如用户是否在线、最近一次加购时间HBase 是全量画像的主存储每一行代表一个用户列族按标签分组Elasticsearch 用来支持运营后台的组合条件筛选比如「近 30 天消费金额大于 1000 且近 1 小时加购过」。三个存储的写入链条通常是通过 Flink 的侧输出分流来实现的主输出写 HBase旁路输出分别写 Redis 和 Elasticsearch。2.3 时效分级不是所有标签都值得实时计算实时计算是有成本的每一个标签都上实时链路会造成状态膨胀和资源浪费。我一般会把标签按时效分成三级T0 秒级标签、T15 分钟分钟级标签、T1 离线回流标签。秒级标签只覆盖「当前会话行为」例如加购状态、正在浏览的商品类目分钟级是常规的窗口聚合结果比如近 1 小时 GMV、近 15 分钟浏览量离线回流则是重计算算法标签的兜底路径。落到 Flink 作业设计上秒级标签用 Processing Time 处理分钟级标签统一走 Event Time 加 Watermark离线回流依赖每日批任务将 Hive 结果写回 HBase。三个时效链路互不干扰每一层都有独立的作业编号和重启策略。这个分级设计要在 Flink 开发前就确定下来因为你一旦把所有标签都塞进同一个流任务后续调优和排障会非常痛苦。3. 用 Flink 搭出实时标签计算管道3.1 数据接入Kafka 消息的解析与分流实时用户画像的源头数据几乎都是埋点日志和业务消息统一进 Kafka 是最常见的架构。Flink 作业从 Kafka 消费 JSON 格式的消息先做规范性校验再根据事件类型分流。事件类型至少包含「曝光」「点击」「加购」「下单」「支付」五类不同事件对应的画像标签完全不同混在一个流里会让窗口聚合逻辑变得很乱。下面是消费入口的参考代码DataStreamString rawStream env.addSource( new FlinkKafkaConsumer( user_behavior, new SimpleStringSchema(), kafkaProps ) ); DataStreamJSONObject validStream rawStream .flatMap(new RichFlatMapFunctionString, JSONObject() { Override public void flatMap(String value, CollectorJSONObject out) { try { JSONObject event JSONObject.parseObject(value); if (event.getString(userId) null) return; if (event.getString(eventType) null) return; event.put(ts, parseEventTime(event)); out.collect(event); } catch (Exception e) { // 脏数据丢入侧输出流后续人工排查 getRuntimeContext().getMetricGroup() .counter(parse_error).inc(); } } }); DataStreamJSONObject sideOutput validStream.getSideOutput(parseErrorTag);这段代码做的事情是从 Kafka 读入行为日志剔除掉缺少 userId 或 eventType 的脏数据并把解析失败的情况单独计数。parseEventTime方法用来从事件内容中提取业务时间戳后续窗口会用它作为事件时间的依据。注意这里不要直接对原始字符串做任何业务处理因为 Kafka 的消息格式很可能在后续迭代中调整规范解析层保持单一职责。对于主流 JSON 解析方案的选型如果数据量不大可以直接用 fastjson 或 Gson但处理吞吐每秒几十万条时建议提前把消息转成 Avro 格式并用 Schema Registry 管理版本。在系统设计阶段就要把「消息中间层是 Kafka Schema Registry」写进架构文档否则后期字段变更会导致 Flink 作业频繁重启而且很难排查是数据问题还是代码问题。3.2 双流 Join行为流与维表流合成标签用户行为数据如果没有商品维度信息支撑很多标签算不出来。例如用户浏览了一个商品你得知道这个商品属于哪个一级类目、什么品牌才能累加类目偏好。这道工序在离线链路里是一次 SQL Join在实时链路里却是大麻烦。实时场景的常见做法是双流 JoinKafka 里除了行为流再建一个商品维度变更流Flink 里用 interval join 按 key 关联两种流。商品维度数据量不大但更新频繁改价、上架、下架直接用广播状态广播维度数据会更简单。下面是广播维表的处理方式MapStateDescriptorString, ProductInfo productStateDesc new MapStateDescriptor(product_info, String.class, ProductInfo.class); DataStreamProductInfo productStream env.addSource( new FlinkKafkaConsumer(dim_product, ...) ); BroadcastStreamProductInfo broadcastProduct productStream .keyBy(ProductInfo::getProductId) .broadcast(productStateDesc); DataStreamTaggedEvent taggedStream behaviorStream .connect(broadcastProduct) .process(new BroadcastProcessFunctionJSONObject, ProductInfo, TaggedEvent() { Override public void processElement(JSONObject value, ReadOnlyContext ctx, CollectorTaggedEvent out) { String productId value.getString(productId); ProductInfo product ctx.getBroadcastState(productStateDesc).get(productId); if (product ! null) { value.put(categoryId, product.getCategoryId()); value.put(brandId, product.getBrandId()); out.collect(new TaggedEvent(value)); } } Override public void processBroadcastElement(ProductInfo value, Context ctx, CollectorTaggedEvent out) { ctx.getBroadcastState(productStateDesc).put(value.getProductId(), value); } });广播维表的核心逻辑是维度流作为广播流每一个行为事件在处理时直接从广播状态中查找商品信息。这种方式避免了每个算子都去外部存储查维度表减少了 Redis 或 MySQL 的压力也规避了异步 IO 的复杂度。需要考虑的问题是广播状态默认存储在内存中商品数量如果达到百万级别每台 TaskManager 都要保存一份全量数据内存很可能成为瓶颈。如果出现这种情况可以将维度数据拆分为高频和低频两部分高频热点商品继续用广播方式低频商品改成查外部维表并加本地缓存缓存失效时间设置在 1 分钟左右既保证实时性又不至于把存储压垮。3.3 会话切割识别用户的一次完整购买旅程很多电商标签以「会话」为粒度例如「本次会话是否产生了加购行为」「平均会话时长」。会话切割在离线 SQL 里可以按用户相邻两条行为时间差大于 30 分钟为界在线处理则需要使用 Flink 的 session window。给会话加标签的代码示例SingleOutputStreamOperatorSessionAggResult sessionAgg taggedStream .keyBy(event - event.getString(userId)) .window(EventTimeSessionWindows.withGap(Time.minutes(30))) .aggregate(new SessionAggregateFunction(), new SessionWindowResultFunction()); // 会话内页面浏览深度 public static class SessionAggregateFunction implements AggregateFunctionTaggedEvent, SessionState, SessionState { Override public SessionState createAccumulator() { SessionState state new SessionState(); state.setPageViewCount(0); state.setHasAddCart(false); return state; } Override public SessionState add(TaggedEvent value, SessionState acc) { acc.setPageViewCount(acc.getPageViewCount() 1); if (addCart.equals(value.getString(eventType))) { acc.setHasAddCart(true); } acc.setBrowseCategory(value.getString(categoryId)); return acc; } }session window 的 gap 参数直接决定会话切割的粒度电商场景一般取 30 分钟业务上可以解释为「用户静默 30 分钟视为一次会话结束」。要注意的是 session window 会触发大量定时器如果用户量级很大每个用户同时存在多个未闭合会话状态量就会成倍增长。解决办法是给会话状态加 TTL例如设置状态 TTL 为 2 小时超过 2 小时的会话强制过期避免内存泄漏。窗口计算完成后会话结果要写入 HBase 的画像宽表同时按用户粒度做实时汇总。会话级标签的一个重要作用是修正统计型标签的偏差用户在 23:59 加购和次日 00:01 加购按自然日窗口它们被分到了两天但按会话窗口它们属于同一次购买意图。如果没有会话维度基于自然日的统计标签很可能对用户行为产生误判这一点在做大促分析时尤其明显。4. 三类画像标签的 Flink 计算实现4.1 统计型标签滑动窗口的几个必调参数统计型标签依赖 Time Window常用的实现是滑动窗口例如「近 1 小时加购次数」「近 7 天支付金额」。Flink 滑动窗口有两个参数一个是 size窗口长度一个是 slide滑动步长。size 决定了标签统计的时间跨度slide 决定了标签更新的频率。这里特别提醒size 和 slide 的比值决定了一个事件会被复制进几个窗口slide 越小窗口重叠越大计算开销也越高。一个比较稳妥的做法是小时级别标签设window size1h, slide5min分钟级标签设window size15min, slide1min。代码示例DataStreamUserBehaviorCount hotTab keyedStream .filter(event - addCart.equals(event.getString(eventType))) .window(SlidingEventTimeWindows.of( Time.minutes(60), Time.minutes(5))) .aggregate(new CountAggregate(), new WindowResultFunction(cartAddCnt)); // 单独统计支付金额双窗口不混用 DataStreamUserBehaviorCount gmvTab keyedStream .filter(event - pay.equals(event.getString(eventType))) .window(SlidingEventTimeWindows.of( Time.hours(24), Time.minutes(10))) .aggregate(new SumAggregate(), new WindowResultFunction(payGmv));每个标签单独开窗口独立计算不要试图在同一个窗口内做多类事件的复杂逻辑判断因为窗口聚合只对单一种类事件有意义。另外窗口函数输出的结果要带上 window_start 和 window_end 字段方便下游存储层做覆盖写和清理也会让依赖画像做报表的同事更好理解数据口径。4.2 规则型标签用状态和定时器拼出判断模型规则型标签的关键在于跨事件判断。拿「高活跃但近 3 天未下单」举例「高活跃」来自窗口计算「未下单」却需要知道用户最近一次支付时间两个信息在不同的时间尺度上必须用 Keyed State 保存长期信息。Flink 代码实现思路如下public class LapseUserFunction extends KeyedProcessFunctionString, TaggedEvent, String { private ValueStateLong lastOrderTime; private ValueStateBoolean isHighActive; Override public void processElement(TaggedEvent value, Context ctx, CollectorString out) throws Exception { if (pay.equals(value.getString(eventType))) { lastOrderTime.update(value.getLong(ts)); } isHighActive.update(computeActive(value)); long currentTs ctx.timerService().currentProcessingTime(); ctx.timerService().registerProcessingTimeTimer(currentTs 3 * 24 * 60 * 60 * 1000L); } Override public void onTimer(long timestamp, OnTimerContext ctx, CollectorString out) throws Exception { Long lastOrder lastOrderTime.value(); if (lastOrder ! null ctx.timerService().currentProcessingTime() - lastOrder 3 * 24 * 3600 * 1000L) { String userId ctx.getCurrentKey(); out.collect(user: userId :risk_lapse); } } }这段代码注册了一个 3 天的 timer每次来新事件都会重新注册。注意定时器的数量要控制好每个用户都有独立的定时器时作业运行一段时间后会导致定时器数量暴涨。优化做法是只对「高活跃」状态为真的用户注册定时器普通用户不必产生无效 timer。4.3 算法型标签Flink 调用外部模型服务的边界算法型标签不在 Flink 内推理原因前面提过。但 Flink 可以负责特征拼接和打分结果落地。假设算法团队已经训练了一个「购买意向模型」输入特征是最近 1 小时浏览次数、加购次数、平均停留时长、用户历史等级Flink 算出这些实时特征后通过异步 IO 调用模型服务获取打分。DataStreamScoreResult scored featureStream .keyBy(event - event.getFlink(userId)) .map(new RichMapFunctionUserFeatures, UserFeatures() { private transient AsyncClient client; Override public void open(Configuration parameters) { client new AsyncClient(http://model-service:8080/predict); } Override public UserFeatures map(UserFeatures value) { value.setScore(client.predict(value.toJson())); return value; } });真正在生产里常用的是 Async I/O异步 IO上面的代码是简化写法实际会通过AsyncDataStream.unorderedWait包装。设置异步等待超时参数一般不能超过 1 秒否则会阻塞下游。算法型标签有一个天然风险模型版本更新会导致同一用户的得分前后不一致所以画像表里要冗余存储打分模型的版本号字段方便问题回溯。5. 画像的存储层设计从 Flink 到 HBase / Redis / Elasticsearch5.1 HBase 宽表Rowkey 与列族设计画像系统的主存储建议选择 HBase原因很简单——写入路径是流式随机写查询路径是 point getHBase 在这两种场景下都有成熟的表现。Rowkey 的设计直接影响热点分布。如果直接把 userId 当 Rowkey数字递增的 userid 会导致写入集中在少数 Region 上。常见处理方法是加盐// 加盐 Rowkey 设计取 userid hash 的前两位拼在原有 userid 前 public static String buildRowKey(String userId) { int salt Math.abs(userId.hashCode() % 100); return String.format(%02d_%s, salt, userId); }在 HBase 里写入的列族可以按标签主题拆分例如行为偏好、交易价值、用户属性、营销敏感度四个列族。每个列族下设多个列列名就是标签编码。HBase 不擅长扫全表所以人群圈选查询不能直接压到 HBase 上而要通过 Elasticsearch 索引找到符合条件的 userId 集合再回表查 HBase 获取详情。写入 HBase 的链路通常直接用 Flink HBase Connector注意要配上批量写入的参数Configuration hbaseConfig HBaseConfiguration.create(); hbaseConfig.set(hbase.zookeeper.quorum, zk1,zk2,zk3); hbaseConfig.set(hbase.zookeeper.property.clientPort, 2181); HBaseSinkString hbaseSink new HBaseSink( user_profile, hbaseConfig, (userProfile, mutation) - { Put put new Put(Bytes.toBytes(buildRowKey(userProfile.getUserId()))); put.addColumn(Bytes.toBytes(behavior), Bytes.toBytes(cart_cnt), Bytes.toBytes(userProfile.getCartCount())); mutation.add(put); } );批量写入参数里最重要的两个配置项是hbase.client.write.buffer和 flush 间隔建议 buffer 调到 2MB 以上否则每条记录都独立提交RegionServer 的写入压力会很大。HBase 的 compaction 会对写入产生毛刺大促前要提前做 major compaction 并关掉自动 split避免在流量高峰触发 Region 分裂。5.2 Redis 里只存短时态数据别当全量存储用Redis 在很多实时系统里被滥用以为 Redis 快就什么都往里面放。用户画像场景下 Redis 的正确落地是保存秒级变化的短时态数据存储结构上建议用 Hash 保存单个用户的一组实时标签HSET user:live:100234 active 1 last_cart_ts 1690000000 page_category electronics EXPIRE user:live:100234 900这个命令里 key 的存活时间只有 15 分钟过期后自动清除确保 Redis 中保存的信息一定是「真实的当前状态」。全量画像不要进 Redis一是内存成本过高二是 Redis 的持久化机制在故障恢复时没法保证数据完全一致画像数据丢了不可接受。删除策略方面Flink 算完新的标签值后可以发送一条 DEL 命令给旧 key 再写入新值避免 Hash key 内字段过多导致读取变慢。如果公司已经切换到云上的 Redis 集群那么主从切换和持久化策略云厂商已经兜底业务侧重点只需要关注内存淘汰策略设置成 allkeys-lru 还是 volatile-lru。5.3 Elasticsearch 支撑运营人群筛选运营同学使用画像系统的频率最高的操作是组合条件圈人找出「近 7 天加购超过 3 次且从未购买过某品牌」的人群。这种查询用 HBase 或者 Redis 写起来效率极低交给 Elasticsearch 则很自然。ES 索引的结构设计建议文档 id 直接用 userId每个标签字段用对应的 ES 类型标签字段ES 类型说明userIdkeyword精确匹配cartCnt7dinteger范围查询lastPayTimedate时间范围筛选preferCategorykeyword多值匹配isHighActiveboolean布尔过滤写入 ES 的速度瓶颈一般在 bulk 大小的设置上Flink Elasticsearch Connector 里bulk.flush.max.actions设置为 1000bulk.flush.max.size.mb设置为 5MBbulk.flush.interval.ms设置为 2000三个参数配合基本能满足分钟级更新要求。如果数据量超过 ES 单集群能承受的极限则要给索引按月滚动保留最近两个月的索引在线更早的数据通过重建索引存入冷节点。6. 端到端验证从数据回溯到 Flink SQL 血排查错画像系统上线不能只看 Flink 作业是否在运行要看最终标签值的准确性和时效性。一个扎实的验证方案分三步走第一步做离线对比抽样第二步按业务事件模拟实时写入观察标签更新时间第三步用 Flink 的 Checkpoint 和指标监控确认链路稳定性。离线对比抽样的做法很简单在 HBase 中随机抽 1 万个用户把实时算出的标签值和离线 Hive 计算的结果对拍。对拍不能要求数值完全一样因为实时计算的窗口边界和离线 SQL 的统计口径不可能完全一致差异在 5% 以内就算通过。如果差异超过阈值优先排查事件时间归属问题和迟到数据处理策略。Flink SQL 在排错阶段很好用。画像系统里不建议所有作业都用 Flink SQL 写但排查 Kafka 某个 topic 的数据格式或者做临时关联查询时Flink SQL 比写 Java 代码快得多。例如验证双流 join 的效果SELECT b.user_id, b.event_type, p.category_name FROM behavior_stream b LEFT JOIN product_dim FOR SYSTEM_TIME AS OF b.proc_time AS p ON b.product_id p.product_id WHERE b.ts CURRENT_TIMESTAMP - INTERVAL 10 MINUTE;数据血缘方面如果是用 Flink SQL 或基于 Catalog 管理的作业可以在平台层记录 Lineage 信息纯 DataStream 作业则靠自定义的 sink 和 source 名称在监控面板里追踪链路。排错的关键是要拿到「事件从进入 Kafka 到写入 HBase 的完整链路时延」这需要在每个算子出口统一打点记录延迟否则某个节点堆积了数据无法快速定位。最后是几个实操里的坑都是反复遇到过的。第一Flink JDBC 连接器连接 HBase 或 MySQL 时默认连接池和重试策略在超时后会直接抛异常要配上连接失败重试和退避间隔否则上游 Kafka offset 不提交导致数据重复消费。第二Watermark 的设置必须结合各业务线数据的最大乱序时间电商晚上大促时埋点上报会比平时延迟更高watermark 设得过于激进会丢掉大量正常事件设得过于保守又会让标签时效性大打折扣推荐先压测再上线。第三状态后端的 RocksDB 开启增量检查点后要监控本地磁盘的使用量状态增长速度快于预期时优先查是否 KeyBy 的 key 粒度过粗或者窗口没清理干净。本文还有配套的精品资源点击获取