
简介《基于Spark的实时用户画像分析系统》是一份技术分享型PDF面向大数据工程师、数据分析师及推荐系统开发者讲解如何利用Spark构建实时用户画像平台解决精准营销与个性化推荐中的用户行为理解问题。资源包为1个PDF文件、2.74MB内容覆盖系统架构、技术栈、功能特点、性能优化与应用场景。文档以优酷用户画像系统为实例重点展示高效筛选器与Join模型优化包括Code Generator、内存列式存储、Bitmap压缩等手段在3~10亿用户、500G数据量、50多个画像维度、5000多个标签的规模下将筛选响应时间控制在2秒群体合并10~20秒对比分析15~20秒。这些真实Benchmark数据对设计高性能实时画像系统极具参考价值。目前已有325人学习下载。适合需要了解用户画像系统整体设计、Spark实时计算实践以及大数据交互式分析优化思路的读者兼具框架讲解与落地细节。1. 为什么等明天的画像让业务慢半拍Spark实时画像系统的定位与第一道选型题用户画像分析系统大家都不陌生——你打开App看到的推荐、营销活动选人、风控规则命中背后都是画像标签在干活。传统做法是T1离线跑批昨天产生的行为今天才落到标签上。这套逻辑在日活几百万、业务决策按天走的时代没问题但到了实时调度、实时营销、实时风控这些场景等一天就是等竞争对手把用户抢走。我下面要拆解的这套基于Spark的实时用户画像分析系统核心思路是跳出离线批处理思维用Spark Structured Streaming消费行为事件流在秒级到分钟级延迟内产出标签并写进特征存储供在线服务查询。适合谁看正在做实时数仓、实时特征服务或者被离线画像延迟坑过的数据平台工程师这篇是把这套系统讲清楚、能落地的一篇实战笔记。2. 实时画像架构选型与标签分层Lambda、Kappa与流批一体怎么选标签体系怎么设计实时画像要解决的第一件事不是写代码而是把架构和标签模型定下来。很多团队上来就写Structured Streaming写到一半发现离线链路和实时链路口径对不上标签重复建设最后越改越乱。动手前把下面两个问题想清楚后面能少返工一半。2.1 实时画像不是把离线任务加速跑从Lambda到Kappa的架构取舍先明确一个观点实时画像不是把离线Spark任务缩短运行周期而是换一条数据流转路径。离线任务每天凌晨读Hive算完后覆盖前一天的全量标签实时任务则持续消费Kafka里的行为事件在事件发生的当下或几十秒内完成统计。两条链路的代码、计算模型、数据源形态完全不同。当前主流的架构选择有三条路我整理成一张对比表架构方案数据链路优点缺点适用场景Lambda离线批处理 实时流处理两套代码并行离线结果稳定可靠实时链路机动灵活两套代码维护成本高标签口径容易不一致公司已有成熟离线数仓实时只是补充Kappa只用流处理历史数据从Kafka重放一套代码、一套口径Kafka需要长期保存全量事件状态管理复杂事件量可控、标签完全由事件驱动流批一体同一套SQL/API跑批和流口径天然统一对引擎版本和SQL兼容性要求高从离线到实时的平滑过渡我见过的大部分公司最终都落在简化版Lambda上离线标签继续保留做兜底和对账实时标签只覆盖高时效场景比如实时营销触达、实时风控、实时调度。Kappa架构在公司里推行阻力很大原因很现实——Kafka保留全量事件几个月的存储成本比想象中高得多而且状态管理出问题时排查难度不小。选型时还有两个前置条件容易被忽略。一是Spark集群本身的资源水位实时任务和离线任务共用集群如果晚上离线大任务把资源吃满实时任务的延迟就会抖动二是埋点数据质量实时画像的底料是事件流如果应用侧埋点缺失、重复上报实时链路会比离线更容易把脏数据暴露给下游。常见的做法是实时任务跑在独立队列或者至少给实时任务打上单独的调度标签。2.2 标签体系怎么设计事实标签、规则标签与模型标签的分层逻辑把架构定了下一步是设计标签体系。我习惯把标签分成三层事实标签、规则标签、模型标签。这三层对应不同的计算频率和数据来源不能混在一个Spark任务里算。事实标签是最底层直接来自行为事件。比如最近5分钟浏览商品数今日下单总金额这类标签用Structured Streaming的窗口聚合就能算出来数据源单一逻辑简单实时性最强。规则标签是在事实标签上做业务规则判断比如近10分钟内加购超过3次的用户标记为高意向用户这类标签需要把事实标签组合起来状态管理更重但规则本身经常被运营调整。模型标签靠机器学习打分比如用户未来24小时购买概率这类标签依赖模型推理一般在离线程算出基础分再在实时链路里根据最新行为做增量修正。设计时有三条原则我踩过坑后才真正理解。第一事实标签要尽量原子化一个标签只表达一个事实便于上层复用第二规则标签的规则要配置化运营改规则时不应该让开发改代码重新发布第三模型标签不要试图完全实时化异步更新往往性价比更高。拿一个农产品行情App举例用户今天搜了三次西兰花是事实标签生鲜价格敏感用户是规则标签明天会再次打开App的概率是模型标签三者时效要求完全不同。这也解释了为什么实时数仓开发工作内容里很大一部分精力花在标签分层和口径治理上而不是写聚合逻辑。标签分层没想清楚后面做规则配置和特征服务时会不断面临这个标签该谁算、算完存哪、多久过期的重复讨论。记住一点实时画像系统里最值钱的不是流处理技术而是标签模型的复用能力。3. 用Spark Structured Streaming搭建实时画像管道Kafka接入、窗口聚合与状态管理参数架构和标签模型定好后就可以动手写管道了。这一章我按三步走先把Kafka事件读进来并解析JSON再做窗口聚合算出事实标签最后在输出阶段用Spark SQL算规则标签。全程用PySpark实现Spark集群搭建好之后这套脚本可以直接在yarn集群上提交运行。3.1 Spark集群就绪后的第一件事Kafka接入与JSON解析的PySpark脚本大多数公司的事件流统一进Kafka实时画像的第一步就是把Kafka里的用户行为事件读成DataFrame。这里有个容易忽略的点Kafka里的value是二进制必须先用from_json解析而且schema要提前和埋点侧对齐。from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, from_unixtime from pyspark.sql.types import (StructType, StructField, StringType, LongType, DoubleType) # 构建SparkSessionshuffle分区数默认200流任务建议显式调小 spark SparkSession.builder \ .appName(realtime-user-profile) \ .config(spark.sql.shuffle.partitions, 8) \ .config(spark.sql.streaming.schemaInference, true) \ .getOrCreate() spark.sparkContext.setLogLevel(WARN) # 事件schema必须与埋点字段严格对齐新增字段要按兼容规则演进 event_schema StructType([ StructField(user_id, StringType()), StructField(event_type, StringType()), # view / click / buy / cart StructField(item_id, StringType()), StructField(category_id, StringType()), StructField(price, DoubleType()), StructField(ts, LongType()) # 毫秒时间戳来自服务端 ]) # 消费Kafka事件流startingOffsets用latest保证只取新事件 raw spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, kafka-1:9092,kafka-2:9092) \ .option(subscribe, user_behavior_events) \ .option(startingOffsets, latest) \ .load() # 先from_json解析成结构体再展开到顶层列 events raw.select( from_json(col(value).cast(string), event_schema).alias(event) ).select(event.*) \ .withColumn(event_time, from_unixtime((col(ts) / 1000).cast(long)).cast(timestamp))这段代码里值得注意的参数有三个。spark.sql.shuffle.partitions默认200离线任务无所谓但流任务每个微批都会触发shuffle200个分区会拖慢小数据量的处理我一般在8到16之间调。startingOffsets用latest意味着任务启动后只消费新事件如果业务上需要补算重启前的数据改成earliest。from_unixtime((col(ts) / 1000).cast(long))这行是把毫秒时间戳转成秒再生成TimestampType直接除以1000会被推断成DoubleType时间就会算错这是读取JSON时最容易翻车的一个细节。3.2 窗口聚合与状态管理30秒窗口内行为标签的计算代码与参数事件读进来之后要做的是按用户维度做窗口聚合。Structured Streaming里的窗口聚合是有状态的系统会为每个窗口保存中间结果watermark用来决定状态什么时候可以清理。from pyspark.sql.functions import window, count, sum, when # watermark设60秒允许事件晚到60秒内仍参与计算 events_with_watermark events \ .withWatermark(event_time, 60 seconds) # 每30秒滑动一次统计每个用户在最近60秒的浏览/点击/加购/购买行为 behavior_stats events_with_watermark \ .groupBy( col(user_id), window(col(event_time), 60 seconds, 30 seconds) ) \ .agg( count(when(col(event_type) view, 1)).alias(view_cnt), count(when(col(event_type) click, 1)).alias(click_cnt), count(when(col(event_type) buy, 1)).alias(buy_cnt), sum(when(col(event_type) buy, col(price))).alias(buy_amount) ) \ .select( col(user_id), col(window.start).alias(window_start), col(window.end).alias(window_end), col(view_cnt), col(click_cnt), col(buy_cnt), col(buy_amount) )窗口参数.window(..., 60 seconds, 30 seconds)第一个参数是窗口长度第二个是滑动步长。每30秒算一次最近60秒的聚合意味着同一个事件会落在两个窗口里好处是标签更平滑坏处是计算量翻倍。如果业务上不需要平滑直接把两个参数都设成30秒就是滚动窗口每30秒只算一个窗口。watermark设60秒我建议按事件的最大乱序程度来定不是拍脑袋。3.3 规则标签用Spark SQL维护把运营规则写成SQL而不是Java代码事实标签算出来后下一步是叠加规则标签。这里的关键是规则不能写在Java/Scala/Python代码里否则运营每次调规则都要改代码重新提交Streaming任务。我一般把事实标签注册成临时视图然后用Spark SQL写规则。def apply_rule_tags_to_redis(batch_df, batch_id): 每个微批先算规则标签再和事实标签一起写入Redis if batch_df.isEmpty(): return batch_df.createOrReplaceTempView(behavior_stats) rule_tags spark.sql( SELECT user_id, high_intent AS tag_name, CASE WHEN buy_cnt 3 OR click_cnt 10 THEN 1 ELSE 0 END AS tag_value, window_end FROM behavior_stats ) # 这里把rule_tags与事实标签合并后写Redis下一章展开用Spark SQL管理规则的好处是规则的可读性大幅提升运营同学能直接看懂买过3次或点过10次就是高意向用户而不是猜代码逻辑。而且SQL天然支持版本管理把SQL文本存进Git每次改动都有历史记录。规则标签的实时性取决于窗口聚合的粒度实际业务里高意向的判定窗口往往是5到15分钟60秒窗口只是示例参数要根据场景调。4. 实时特征落库与在线查询Redis在画像服务里的数据结构选型与写入参数Spark把标签算出来后如果没有一个在线存储让业务服务毫秒级拿到标签这套系统就白做了。实时特征服务的核心诉求是低延迟、高并发、支持过期——这三个词基本把存储选型指向了Redis。4.1 为什么Redis是实时画像的默认在线存储对比HBase与ClickHouse的点查能力先看一张选型对比这是我被问得最多的一个问题存储延迟数据模型过期机制运维成本点查场景Redis毫秒级Hash/Set/ZSet等丰富结构原生TTL低内存成本高按user_id取全部标签HBase毫秒级KV列族无原生TTL需配置高依赖HDFS按user_id取指定列ClickHouse亚秒级列式表无靠分区淘汰中不适合高频点查用户画像的点查模式是按user_id取一批标签Redis的Hash结构完美匹配这个场景key是profile:user:{user_id}field是标签名value是标签值。一次hgetall就能拿到一个用户的全部实时标签RTT只有一次。HBase在数据量上更强能存PB级画像数据但一个在线实时画像系统如果只靠点查几百GB内存的Redis集群足够覆盖几亿活跃用户的高频标签。ClickHouse更适合分析师跑画像人群SQL不适合在线服务调用。选Redis还有一个隐性的便利TTL机制天然契合实时标签的时效。离线画像标签一存就是一个月实时画像标签可能30秒后就过期了。Redis的过期策略让数据自动淘汰不需要额外写清理任务这在实时数仓开发里是一个很大的省心事。4.2 标签写入与毫秒级读取foreachBatch加Pipeline的落地代码与连接池参数具体写入我推荐用foreachBatch而不是foreach。foreachBatch每个微批调用一次可以复用Redis连接、合并写入、在批内再聚合foreach是每条记录调用一次只适合写入量极小的场景。下面是完整的写入代码。from redis import Redis # decode_responsesTrue 让hget返回str而不是bytes省去手动解码 r Redis(hostredis-feature, port6379, db0, decode_responsesTrue, socket_connect_timeout3) def write_tags_to_redis(batch_df, batch_id): 每个微批将事实标签和规则标签批量写入Redis if batch_df.isEmpty(): return batch_df.createOrReplaceTempView(behavior_stats) rule_tags spark.sql( SELECT user_id, high_intent AS tag_name, CASE WHEN buy_cnt 3 OR click_cnt 10 THEN 1 ELSE 0 END AS tag_value, window_end FROM behavior_stats ) # 用pipeline批量发送减少网络RTTtransaction关闭避免MULTI/EXEC开销 pipe r.pipeline(transactionFalse) for row in rule_tags.collect(): key fprofile:user:{row.user_id} pipe.hset(key, mapping{ view_cnt_60s: 1, # 实际应从事实标签透传 high_intent: row.tag_value, updated_at: row.window_end.strftime(%Y-%m-%d %H:%M:%S) }) pipe.expire(key, 60) # 标签60秒过期与窗口时长保持一致 pipe.execute() query behavior_stats.writeStream \ .foreachBatch(write_tags_to_redis) \ .outputMode(update) \ .trigger(processingTime30 seconds) \ .option(checkpointLocation, hdfs:///user/spark/checkpoint/profile-v1) \ .start() query.awaitTermination()这段代码要重点理解三个参数。pipeline(transactionFalse)把一批命令打包发送命令数多时吞吐能提升好几倍关掉事务避免MULTI/EXEC的额外开销。expire(key, 60)让标签在60秒后自动过期下游服务读不到过期标签时再回退到离线兜底值这是实时画像和离线画像衔接的好设计。checkpointLocation必须放在HDFS或S3这类共享存储上如果放在本地磁盘任务重启后状态就丢了这个坑下面具体说。在线服务读取侧的代码简单得多但有一个参数值得注意要用连接池不能在每次请求里新建连接。import redis from redis.connection import ConnectionPool # 在线服务用连接池复用连接max_connections按峰值QPS估 pool ConnectionPool(hostredis-feature, port6379, db0, decode_responsesTrue, max_connections50) r redis.Redis(connection_poolpool) # 一次RTT拿到用户全部实时标签 tags r.hgetall(profile:user:9527)读取侧真正的坑在decode_responses。如果读写两侧有一侧没开这个参数hgetall会返回一堆bytes对象到代码里还得逐个decode小问题但很烦。另外给Redis配置maxmemory-policy allkeys-lru能保证内存写满时优先淘汰最久没访问的标签这个策略对画像这种热标签反复读、冷标签被淘汰的场景非常合适这是我用Redis存画像标签时最常调的一个参数。5. 实时画像的5个坑与排查乱序、状态膨胀、热点倾斜、checkpoint、Redis背压实时画像和离线画像最大的差别在于离线任务跑挂了可以重跑实时任务跑挂了会丢数据、会延迟。这一章我用几个真实场景梳理最常遇到的坑每个问题按现象、原因、解决的顺序讲排查思路可以直接套用。5.1 事件乱序让标签算错Watermark设60秒还是5分钟现象用户行为标签在某些时段异常偏低比如晚间App端上报的事件明显少了但离线日志显示用户行为很活跃。排查发现是用户设备网络状态差异导致事件到达Kafka的时间不统一4G弱网下的事件可能延迟几分钟才到而窗口聚合早就关闭了。原因Structured Streaming按事件时间做窗口但默认情况下事件晚到会被直接丢弃watermark决定了允许晚到多久。watermark设太小晚到事件全部丢失设太大状态保存时间拉长内存压力变大。解决watermark的设置不能凭感觉我一般先看事件延迟分布。用Kafka consumer的records-lag-max和埋点里的客户端时间戳做对比统计P95延迟。如果95%的事件在90秒内到达watermark设120秒到180秒比较稳。多丢一点数据还是多撑一会内存是个权衡纯实时营销场景我倾向设大一点因为标签算错比延迟更伤害业务对账场景才需要严格控制watermark。5.2 状态膨胀拖垮executor状态清理与Spark内存配置的血泪经验现象任务刚上线时一切正常跑了一周后executor频繁Full GC任务处理延迟从30秒恶化到5分钟最终OOM。原因Structured Streaming的状态存储会为每个窗口、每个key保存聚合中间值。窗口长度大、key数量多、watermark设得久状态就会持续膨胀。尤其是用户规模大但窗口又不合理的任务状态增长速度远超预期。解决三个手段要一起用。第一缩窗口——如果业务允许把窗口从1小时改成15分钟状态量直接降四分之三第二收紧watermark——晚到容忍度从10分钟降到3分钟状态清理会更激进第三给executor配off-heap内存状态存的是Java对象堆内存不够时会触发GC风暴在提交参数里加spark.memory.offHeap.enabledtrue和spark.memory.offHeap.size2g能缓解。还有一个偏方把计算拆成两个任务一个算窗口内的增量一个单独做状态清理见过有人这么干但维护成本高不推荐。5.3 热点用户造成数据倾斜加盐聚合在流任务里的正确姿势现象同一个Streaming任务里大部分executor都空闲只有一个executor跑满CPU。查看Spark UI发现某个task的输入数据量是其他task的上百倍这是一个超级用户产生的行为事件占了大头。原因groupBy(user_id)在shuffle时按user_id哈希到固定分区热点用户的所有事件都落进同一个分区分区内的任务自然成为瓶颈。解决离线批处理里的加盐思路在流处理里要谨慎用。正确的做法是事件写入Kafka时就按user_id % shard_num分片保证同一个用户有序落盘聚合时先按分片聚合再做第二层合并。-- 第一段聚合按分片键和窗口先算出局部结果 SELECT shard_id, user_id, window_start, window_end, COUNT(*) AS cnt FROM events_with_shard GROUP BY shard_id, user_id, window_start, window_end; -- 第二段聚合合并分片结果 SELECT user_id, window_start, window_end, SUM(cnt) AS total_cnt FROM stage1 GROUP BY user_id, window_start, window_end;但也要说清楚这种加盐方式只对普通热点有效如果某个超级用户的量级是普通用户的上百万倍拆到16个shard还是单台executor扛。超大用户的终极方案是单独识别、单独走一条高配置链路也就是热点隔离特征服务里对头部用户单独建桶这是各大厂的标准做法。5.4 checkpoint恢复失败升级之后任务起不来的后悔药现象集群的Spark版本从3.3升级到3.4同时调整了事件schema加上了一个字段。重启Streaming任务后直接报CheckpointMissingException或序列化错误任务起不来。原因Structured Streaming的checkpoint里保存了状态存储的schema和操作符元数据Spark小版本升级或计算逻辑变更都可能导致checkpoint与当前代码不兼容。Structured Streaming官方明确不保证跨版本恢复这个问题遇到过一次就长记性了。解决我的习惯是——每次升级或者改schema之前先把checkpoint目录完整备份一份到另一个路径如果代码逻辑大变直接放弃旧checkpoint用一个新appName和新checkpoint目录重新跑必要时用startingOffsets: earliest从Kafka最早偏移消费历史数据补算。还有一点checkpoint目录一定不能用回收站HDFS的trash回收策略可能在任务没结束就清掉老版本的状态这个坑我踩过一次后给checkpoint目录单独配了禁用trash的权限。5.5 Redis写入成为瓶颈消费延迟排查与背压处理现象Kafka消费延迟持续上涨但看Spark UI发现executor的CPU使用率不高GC也不频繁。查Streaming指标发现每个微批的处理时间越来越长从3秒涨到30秒。原因foreachBatch里的Redis操作拖慢了整个微批。如果没用pipeline而是一条条写Redis每条都要走一次网络RTT一个微批几万条记录光写Redis就要几十秒Kafka端自然堆积。解决第一步把逐条写入改成pipeline批量写入能解决大部分问题第二步关闭transactionTrue因为事务在pipeline里会触发MULTI/EXEC额外的往返让性能回到逐条水平第三步如果写入量实在太大考虑换架构——把Redis写入从流任务里拆出去流任务只写到Kafka sink再用独立的消费者服务通过异步批量写Redis。拆链路会增加维护成本但吞吐瓶颈从Spark转到Redis时这是最通用的解法。6. 画像准不准、值不值延迟监控、抽样比对与AB验证方法实时画像系统上线不等于结束真正的考验是证明标签准、延迟低、业务有价值。这一章给三个验证手段都是我每个新标签上线必做的事。6.1 标签准确性验证和离线画像对账的抽样方法实时标签跑得再快如果算错了下游业务就是在错误数据上做决策。我每个新标签上线前都要先做对账每天凌晨用离线任务重算一遍同样的标签再和实时标签做对比。抽样方法很关键不能只比对活跃用户要按用户分段分层抽样头部用户、中部用户、尾部用户各抽一部分避免都被头部用户的数据淹没。对账的核心指标是标签一致率两个链路算出来的标签值在允许误差范围内的比例。不同的标签类型对账口径不一样事实标签比的是数值误差规则标签比的是命中率模型标签比的是排序相关性。做一个简单的对账表验证维度方法参考阈值事实标签准确性抽样用户离线重算后与实时值比对数值误差小于5%规则标签命中一致率抽样用户比较规则是否同时命中一致率不低于90%端到端延迟在Redis记录写入时间与事件时间差P95小于3分钟数据完整性实时标签用户数与离线UV对比偏差小于10%6.2 实时标签的性价比排序先做哪些标签后做哪些标签实时计算很贵不是所有标签都值得实时化。我给标签做优先级排序时只看一个维度标签晚一小时拿到业务损失是多少。实时营销触达、实时风控拦截这几个场景晚一小时等于钱没花出去或者坏账多了一笔必须实时推荐系统粗排里用户兴趣标签晚15分钟影响可以接受人群洞察、用户分层这类运营分析标签走离线就够了不需要占用实时资源。我自己的习惯是每次上线一个新标签先写对账SQL再写规则表达式标签上线前先在离线表上跑一遍历史数据分布确认口径无误再提交Streaming任务。这套流程多花半天时间但能省下后面排查口径不一致的好几个通宵。实时画像从来不是上了一个流任务就结束的事把验证机制和标签分层做好系统才真正能跑稳、跑久。希望帮到你。本文还有配套的精品资源点击获取