Spark实时用户画像:标签建模、流式计算与集群调优实战
简介基于Spark的实时用户画像分析系统是一份面向大数据工程师、数据平台架构师及推荐/营销算法团队的技术PDF聚焦如何利用Spark生态构建面向海量用户的实时画像平台。内容以优酷用户画像系统为实例完整讲解了系统架构、技术选型、筛选器设计、Join模型优化与存储方案覆盖Spark、Hadoop、Scala、ANTLR、SQL等核心组件并给出3~10亿用户、500G数据、50维度、5000标签下的性能Benchmark能帮助读者理解从数据采集、标签计算到精准投放/推荐落地的全链路方法。资源为单个PDF文件大小2.74MB便于阅读和收藏。目前已有324人学习适合希望系统掌握实时用户画像体系设计思路并借鉴优酷工程实践经验的大数据从业者。1. 实时用户画像的落地形态Spark在这里解决什么问题运营同学早上发来需求给最近30分钟活跃但过去7天没下单的用户发一张优惠券。如果走T1离线数仓等标签算出来用户可能已经进入沉睡期如果走自研流处理又很难复用已有的数据清洗逻辑。实时用户画像系统的本质是让“用户标签”从离线批量更新变成事件驱动更新在秒级或分钟级完成标签写入。Spark在这个场景里不是唯一选择但它有一个很实际的优势同一套DataFrame代码既可以跑批也可以跑流。常见做法是用Spark Structured Streaming消费Kafka事件借助内存计算和分布式执行把ETL逻辑直接复用到实时链路。这个能力对团队很重要因为画像的维度会随时间膨胀离线开发一套逻辑、实时再开发一套维护成本会翻倍。下面的内容面向需要从0搭建或重构实时画像的数据工程师、数据平台工程师和沾数据后端的业务开发。你先能看到标签体系怎么建再看到Spark实时ETL怎么写、集群和内存怎么调最后给出线上质量校验和迭代技巧。全程都有可复现的代码和参数照着落地基本不会跑偏。2. Spark实时用户画像的标签体系与数据建模实时用户画像不仅仅是“存一个用户维度表”它决定你后续所有分析的边界。一个常见的错误是先定技术再补标签导致Spark任务里堆满if else却不知道标签的业务口径。我的习惯是先按业务需求拆标签维度再回到Spark里做特征计算。2.1 从业务需求反推标签维度标签体系可以从更新频率和计算逻辑两个角度拆分。按计算逻辑分最常用的是三类规则型标签直接处理事件字段统计型标签依赖窗口聚合算法型标签来自模型评分。更新频率上还要区分状态快照与累计指标例如“是否在直播间”只要存最后一次状态“过去7天访问次数”则需要逐天累加。两者对实时写入路径的要求完全不同。下面这张表是我在实际项目里整理出的标签维度划分重点看实时入口与更新频率的搭配。标签类别典型示例更新频率实时入口基础属性设备型号、注册渠道、会员等级天级或事件驱动注册/登录事件行为频率近5分钟浏览次数、近1小时加购数分钟级埋点/交易事件状态标签是否在直播间、是否正在售后秒级状态变更事件偏好标签类目偏好、价格带偏好小时级算法评分任务这张表看着简单但实际建模时有一个关键点不同更新频率会影响Spark的存储选型。秒级状态标签要放Redis分钟级统计标签适合放HBase或Doris天级偏好标签可以留在数仓里。后端第3章会具体说明存储对接方式。设计阶段还必须输出标签字典包含标签名、中文名、计算逻辑、更新频率、责任人。没有字典画像跑一个月后没人知道“高活跃用户”到底指什么。2.2 用DataFrame做标签数据的清洗与加工清楚了标签分类下一件事就是把原始日志变成标准客户行为事件。这里我强烈推荐直接用PySpark DataFrame而不是RDD。DataFrame有内置的schema推断和谓词下推写出来的ETL脚本更像SQL不同资历的工程师都能看懂。下面这个最小示例同时演示了Kafka读取、schema解析和规则标签生成。from pyspark.sql import SparkSession from pyspark.sql.functions import to_timestamp, when, col spark SparkSession.builder \ .appName(user_profile_label_etl) \ .enableHiveSupport() \ .getOrCreate() # 读取Kafka中的埋点日志这里以二进制格式为例 raw_df spark.read.format(kafka) \ .option(kafka.bootstrap.servers, kafka001:9092,kafka002:9092) \ .option(subscribe, user_behavior) \ .load() # 解析value并筛选关键字段 event_df raw_df.selectExpr(cast(value as string) as json_str) \ .select( from_json(json_str, userId STRING, eventType STRING, timestamp STRING, itemId STRING) ) \ .select( col(userId), col(eventType), to_timestamp(timestamp, yyyy-MM-dd HH:mm:ss).alias(event_time) ) # 生成规则标签加购未支付、活跃用户、流失风险 labeled_df event_df.withColumn( label, when(col(eventType) add_cart, cart_without_pay) .when(col(eventType) view, active_user) .when(col(eventType) order, buyer) .otherwise(unknown) ) labeled_df.write.format(parquet).mode(overwrite).save(/tmp/user_labels)这个脚本的读流程很典型从Kafka读取原始二进制消息用from_json拆出userId、eventType等字段然后通过withColumn做规则判断。注意from_json需要手写schema在实时链路里很常见因为流式数据不像离线表有统一元数据。实际ETL脚本还会把字段解析、事件清洗、标签计算拆成不同函数便于Spark的物理计划复用。如果字段很多建议用StructType而不是字符串拼schema否则字段一旦嵌套多维护会非常痛苦。2.3 画像存储模型与宽窄表选择标签计算完不是落Parquet就结束了最终要供在线查询。离线数仓里我们习惯用宽表一个用户一行几十个字段。但实时画像场景不适合把全部标签塞进同一张宽表原因有两个一是标签更新频率不同写宽表会导致频繁Update全行二是Spark实时任务要保证低延迟窄表按标签名写更容易做局部更新。我一般会把画像表设计成窄表主键是user_id tag_nametag_value存放JSON或普通类型update_time记录最后更新时间。在Redis里用hash结构存储key是user_idfield是tag_namevalue是tag_value。这一设计在启动Spark实时任务前就定好后面集群调优、存储扩容才不会改模型。窄表映射到Spark的分析路径也很自然。做用户群体分析时可以通过Spark SQL把行转列变成宽表供BI使用在线查询时同一份窄表又能按RowKey直接Get。这样实时链路和离线分析共用一份主模型避免数据口径二次加工。3. 用Structured Streaming把画像更新延迟降到秒级第2章把标签模型和离线ETL脚本理清了但离“实时”还有一步如何让Spark持续消费事件并把计算结果写到在线存储。这里明确说不建议再用老牌的Spark Streaming基于DStream API应该用Structured Streaming。3.1 为什么选Structured Streaming事件时间与状态管理Structured Streaming把流表当无边界的DataFrame你可以用同样的select、groupBy、agg操作来做实时计算。对于实时画像最有价值的是水印和事件时间语义。埋点日志往往因为前端重试或网络抖动乱序如果用处理时间计数会得出“今日活跃下降”这种错误结论。Structured Streaming允许用watermark忽略迟到的数据并在状态存储里保留中间结果。3.2 实时画像ETL核心代码窗口统计与标签合并下面是一个典型的实时画像ETL框架我整理成可以直接改的脚本。这个脚本会消费Kafka中的浏览和加购事件按用户ID做5分钟滑动窗口统计然后将结果upsert到具有主键的MySQL表中。这里用MySQL做示例生产环境可以换成HBase或Redis。from pyspark.sql.functions import window, col, count, when, to_timestamp, from_json stream_df spark.readStream.format(kafka) \ .option(kafka.bootstrap.servers, kafka001:9092) \ .option(subscribe, user_behavior) \ .option(failOnDataLoss, false) \ .load() \ .selectExpr(cast(value as string) as json) events stream_df.select( from_json(col(json), event_schema).alias(e) ).select( col(e.userId).alias(user_id), col(e.eventType).alias(event_type), to_timestamp(col(e.ts), yyyy-MM-dd HH:mm:ss).alias(event_time) ) # 窗口聚合计算行为标签 profile_updates events \ .withWatermark(event_time, 5 minutes) \ .groupBy( window(event_time, 5 minutes, 1 minute), user_id ).agg( count(when(col(event_type) view, 1)).alias(view_count), count(when(col(event_type) add_cart, 1)).alias(cart_count) ) def write_upsert(df, epoch_id): df.write.format(jdbc) \ .option(url, jdbc:mysql://profile-db:3306/profile) \ .option(dbtable, user_realtime_label) \ .option(user, profile) \ .option(password, secret) \ .option(batchsize, 1000) \ .mode(append) \ .save() query profile_updates.writeStream \ .foreachBatch(write_upsert) \ .outputMode(update) \ .trigger(processingTime60 seconds) \ .queryName(user_profile_streaming) \ .checkpointLocation(/data/checkpoint/profile) \ .start() query.awaitTermination()这里需要重点说明几个参数。withWatermark允许事件时间晚到5分钟超过阈值的旧数据会直接丢弃这是控制脏数据进入核心链路的关键。window定义5分钟窗口、1分钟滑动步长保证标签口径是近5分钟。outputMode(update)只输出发生变化的结果适合画像更新这种不断叠加的场景。foreachBatch则是一个万能出口适合JDBC这类不支持流式写入的目标端可以在批内做去重再写。checkpointLocation必须指向HDFS或云盘不能放本地/tmp否则任务重启时会丢offset。另外failOnDataLoss建议设falseKafka提前删除消息时任务不会直接崩溃。3.3 在线存储怎么选Redis还是HBaseSpark实时任务计算完标签后需要写入在线系统。不同业务对读取延迟、更新频率要求不同。在我经历过的项目里选型主要看两点单个用户标签的读写并发以及标签集合的大小。下面这个表格总结了常见选型。存储读取延迟实时更新方式适用场景Redis(Hash/Stream)毫秒级HINCRBY/HSET秒级状态标签、推荐系统特征HBase/Hudi毫秒~几十毫秒Put on same RowKey海量用户、标签数量多、需要冷热分离MySQL(PK upsert)几毫秒~几十毫秒ON DUPLICATE KEY UPDATE内部运营后台、数据量万级以内如果并发量很高且要支持画像分析可以选HBase作为离线结果与实时结果合并的存储层再通过Phoenix或Presto做分析查询。之前设计的窄表模型映射到HBase的RowKey(user_id) Qualifier(tag_name)非常自然。如果只是推荐系统取特征Redis/hash结构更合适但要记得给key设置TTL否则长期用户的状态会撑爆内存。4. Spark集群搭建与内存/ETL调优让实时画像稳得住实时任务不比离线夜批出问题就得立即恢复机器挂了三五分钟都肉眼可见。这一章把从Spark集群搭建到实时任务稳定运行的关键配置讲清楚。标题里的“集群搭建”对应这里不只是启动master/worker还要考虑提交模式与资源隔离。4.1 最小Spark集群搭建与安装使用真实生产环境用YARN/SBP返回调度但很多团队会把“集群搭建”理解成用独立模式跑通然后换YARN反而栽跟头。我建议先用Standalone模式跑通再迁移到YARN快速验证任务代码。下面的步骤以Spark 3.x为例从Apache官网下载spark-3.x.x-bin-hadoop3.2.tgz部署在3台机器上。# 下载并解压到/data/soft tar -zxvf spark-3.x.x-bin-hadoop3.2.tgz -C /data/soft/ ln -s /data/soft/spark-3.x.x-bin-hadoop3.2 /data/soft/spark # 配置workers cd /data/soft/spark/conf cp spark-env.sh.template spark-env.sh echo JAVA_HOME/usr/java/jdk1.8.0_212 spark-env.sh echo SPARK_MASTER_HOSThost01 spark-env.sh echo SPARK_WORKER_CORES8 spark-env.sh echo SPARK_WORKER_MEMORY24g spark-env.sh echo SPARK_WORKER_OPTS-Dspark.worker.cleanup.enabledtrue spark-env.sh # 启动集群 /data/soft/spark/sbin/start-master.sh /data/soft/spark/sbin/start-workers.sh这段命令的核心是SPARK_WORKER_MEMORY和SPARK_WORKER_CORES它们决定这个节点能被多少executor共享。一个常见错误是把本机物理内存全配给worker却不给操作系统和DataNode留余量结果节点卡死。建议保留20%给系统。4.2 实时任务的资源与内存参数设置集群搭建完后重点就变成Spark作业内的参数。实时画像任务采用spark-submit方式提交通常我是这样设置执行参数的配置项建议值解释--executor-memory4-6g单个executor堆内存根据对象大小调整--executor-cores2-3CPU核数太高会严重浪费同时增加GC时间--driver-memory3gdriver端保存状态会占内存别默认给1gspark.sql.shuffle.partitions200~400默认200若数据量小过大只是浪费spark.memory.offHeap.enabledtrue开启堆外内存减少堆内对象拷贝spark.kryoserializer.bufferSize64m使用Kryo序列化避免Java序列化带来的GC压力除了executor memory还要注意spark.executor.memoryOverhead。YARN模式下默认是executor-memory的10%如果任务频频报内存Overhead多是因为Spark内部有大量网络缓冲和直接内存建议手动设到1g以上。另外Structured Streaming的checkpoint目录必须放在可靠的存储比如HDFS不能放本地/tmp。4.3 ETL脚本中的常见性能瓶颈与spark-etl优化实时画像的ETL脚本除了读Kafka还会做维度关联、身份归一化等操作。最容易出现的性能问题是shuffle倾斜个别user_id因为爬虫或者刷单流量巨大导致某个任务节点长时间无法结束。有一组优化手段可以套用。使用随机前缀两阶段聚合处理分钟级热点用户在groupBy前把user_id打散成瞬时key对维表用broadcast join只要维表小于spark.sql.autoBroadcastJoinThreshold默认10MB避免shuffle关闭不必要的Spark UI日志保留策略避免大段时间的driver overhead使用DataFrame API而不是RDD让Spark Catalyst优化器执行列剪裁和谓词下推。在实时画像场景最有效的是分区内限制记录数可以在foreachBatch里先filter再写入防止小文件爆炸。另一个容易忽略的点是如果同一个Spark任务同时跑多个流要让不同的queryName绑定不同的checkpoint路径否则Kafka offset会相互覆盖。5. 实时用户画像系统的质量校验与增量迭代技巧系统上线后不能只盯着延迟还得回答“标签算得准不准”。我这里的做法是用离线Spark任务抽样对账同时对晚到数据做标签回刷。5.1 用离线任务校准实时标签写一个离线任务从数仓取同一时间段的用户行为明细用相同的口径计算标签然后和实时存储的标签做diff。SELECT wr.user_id, wr.view_count AS realtime_cnt, off.view_count AS offline_cnt, IF(wr.view_count ! off.view_count, 1, 0) AS diff_flag FROM (SELECT user_id, view_count from realtime_tag WHERE dt#{predict_date}) wr JOIN (SELECT user_id, count(*) AS view_count from dwd_user_event WHERE event_time between #{start} and #{end} GROUP BY user_id) off ON wr.user_id off.user_id HAVING diff_flag 1 ORDER BY diff_cnt DESC;需要注意实时任务和离线任务取绝对相等意义不大因为流处理窗口是滑动计算离线按天切分会有边界数据。不要直接对diff_flag做全量告警而应该设一个容忍比例比如偏差超过5%才触发。5.2 处理延迟数据与标签回刷Structured Streaming的watermark会丢到5分钟后的数据但这个值并不绝对。某些用户的行为轨迹可能延迟半小时导致画像“卡住”。常见做法有两种一是调大watermark到10-15分钟另一种是写一个离线标签回刷脚本扫描过去一天的事件日志重新计算丢数据较多的标签。我比较常用第二种因为watermark调大会增加状态存储成本。回刷脚本是一个Spark批任务只需要在spark-submit时传入一个时间参数spark-submit --class com.example.RollbackLabel \ --master yarn \ --deploy-mode cluster \ --executor-memory 6g \ --executor-cores 3 \ /app/profile/tag-backfill.jar \ --begin 2025-03-01 --end 2025-03-01 --targetTag view_count这个脚本以“T1”的方式重新跑一遍用户行为然后把计算出来的标签批量覆盖进Redis/HBase这样第二天早晨数据就齐了。5.3 用StreamingQueryListener做延迟监控与告警比起写一堆自定义监控我更建议直接使用Spark StreamingQueryListener每5分钟记录一次进程的inputRowsPerSecond、processedRowsPerSecond、事件时间延迟等指标有异常就直接推送到钉钉/企业微信。核心逻辑是在onQueryProgress回调里读progress对象提取关键字段。下面是一个典型的progress输出{ id: user_profile_streaming, name: user_profile_streaming, timestamp: 2025-03-01T10:00:0008:00, batchId: 15, inputRowsPerSecond: 1000.0, processedRowsPerSecond: 600.0, durationMs: 5000 }在实际工程里我会把回调里拿到的processedRowsPerSecond写入InfluxDB配合Grafana看趋势。如果连续3个微批processedRowsPerSecond低于正常值的50%就触发告警。这个方式比观察Spark UI更及时因为UI最多延迟几十秒而告警可以直接打到值班群。本文还有配套的精品资源点击获取