资讯动态

基于Spark Structured Streaming与Kafka的实时用户画像系统构建指南

发布时间:2026/9/17 20:23:43 来源:尧图企业网站定制
简介一份深入讲解基于Spark的实时用户画像分析系统的技术PDF适合大数据开发、算法工程及推荐系统相关从业者学习参考。内容从系统架构入手覆盖用户画像系统、实时计算引擎、存储系统与交互Server的协作方式并详细介绍了Spark、Hadoop、Scala、Java、ANTLR、SQL等核心技术选型。文档还展示了实时画像分析、精准营销推荐、高效筛选器、大数据量处理等功能特点以及在内存计算、列式存储、Bitmap压缩、Join模型优化等方面的性能调优思路。以优酷真实场景为案例给出了3~10亿用户、500G数据量、50余画像维度、5000余标签下的Benchmark数据对构建生产级用户画像系统有较强参考价值。压缩包内共1个PDF文件整体大小2.74MB已有324人学习。1. 实时用户画像不是“跑批提速”Spark 的切入点在哪里用户画像是数据分析里被讲烂、但实时化之后完全变味的一件事。传统做法是凌晨跑 Hive 或 Spark 批任务T1 更新标签次日运营看到的数据已经是昨天的用户。实时用户画像分析系统要解决的是把链路压到分钟级甚至秒级用户点了一个商品、触发一次搜索、看完一段视频行为在几分钟内就要反映到“近期偏好”“购买意向”“活跃等级”等标签上。Spark 在这里不是最快的流引擎而是适配成本最低的流批一体底座——Structured Streaming 复用 DataFrame API让同一套计算逻辑既能吃实时流也能回放历史做补偿。打算自己搭画像系统的后端和数据工程师这份拆解可以直接当落地参考。2. 实时画像链路设计Kafka 主题、消费语义与标签服务落点2.1 先定主链路Kappa 架构为什么比 Lambda 更适合用户画像Kappa 和 Lambda 的争论在画像场景里其实没有悬念。Lambda 要求同一套标签逻辑在流和批里各写一遍流里用 Structured Streaming批里用 Spark SQL两边的窗口边界、去重规则稍微差一点同一个用户的标签就会开始漂。更麻烦的是等线上发现两组标签对不上你很难判断是流算错了还是批算错了排障成本直接翻倍。Kappa 的思路是数据只留一份全部进 Kafka实时作业从增量位点消费需要重算时就近回放。用户画像标签天然是时间状态昨天的活跃等级不影响今天但近 7 天购买窗口在持续叠加所以 Kafka 的消息保留期一般拉到 7 到 15 天正好覆盖画像计算需要的最大回溯范围。主链路定下来后面的存储选型和补偿机制都围绕它展开不再维护双份口径。2.2 Kafka 主题怎么拆多主题 user_id 做 key主题划分直接影响 Spark 作业的并行度和解析复杂度。单主题 user_events 最省事所有行为统一入流但页面浏览、下单、搜索三类事件的字段差异大schema 只能往宽松了设计后面每个作业都要做清洗。按事件域拆成 page_views、orders、searches 三个主题各主题 schema 紧凑还能独立调并行度——下单量小orders 分 8 个分区就够浏览量大page_views 可以分 64 个分区。方案优点代价适用单主题 user_events接入简单一套 schema 通吃字段松每个作业都要过滤清洗事件量小、人力紧张的初期团队多主题分事件域schema 紧凑并行度可独立调跨主题状态合并要额外处理事件量分域明显的生产环境事件消息体用紧凑 JSONpage_views 的一条样例长这样{user_id:u_10241,event_time:2026-05-14T11:23:0708:00,event_type:page_view,scene:home_feed,item_id:sku_9832,props:{ref:push_camp_071}}两个约定要写进规范。event_time 用带时区的 ISO8601不能传服务端入库时间否则窗口聚合在跨时区时全部错位Kafka 消息 key 必须取 user_id同一用户的行为落在同一分区后续按用户聚合不需要重分区同一用户的乱序消息也尽量按 partition 顺序处理。2.3 选 Spark 而不是 Flink流批一体维护成本才是决定因素真正决定选型的不是单条吞吐那种跑分而是这台作业交到你手里之后一年要花多少维护时间。Structured Streaming 是微批模型一次 trigger 处理一小批数据和 Flink 的逐条处理比延迟不占优势但换来三个实打实的好处。第一API 复用。readStream 和 read 出来的都是 DataFrame聚合、join、UDF 写法完全一致这就是 Spark 的 DataFrame 编程模型直接搬到流上团队里会写批处理的人不用重新学算子。第二端到端 exactly-once。Kafka 读取位点、状态、输出统一交给 checkpoint 管理配合幂等写出作业重启不会产生重复标签。第三watermark 内建晚到数据可以声明容忍窗口不需要自己维护迟到事件的去重。代价也要说清楚状态算子在流里支持有限复杂跨事件的全状态聚合往往要下探到 mapGroupsWithState 手写不如批处理自由微批延迟决定了 trigger 通常设在 1 到 5 秒。画像标签的时效是分钟级完全覆盖。另一个现实因素是运维没有专职流计算团队的公司在 Spark 上多挤一份资源比单独养一套 Flink 集群健壮得多。2.4 标签服务的两种形态查询型与推送型画像系统的出口形态决定了下游怎么对接。常见做法分两类。一类是查询型推荐、运营按 user_id 即查即用Redis 里一个 hash 存一个用户的全部标签HSET 写入、HGETALL 读出单次访问毫秒级。另一类是推送型标签变化幅度超过阈值时把变化事件写回 Kafka 的 profile_updates 主题下游 CRM、push 系统再消费。推送型的判定在 foreachBatch 里完成读旧值、比新值变化量够大才发避免每个批次空转。3. 用 Spark Structured Streaming 搭建实时标签计算的最小链路3.1 spark集群搭建与作业提交三个必须先定的配置假设你有 24 个 executor 可用下面这份 spark-submit 是生产环境的常见起点。spark的安装与使用本身门槛不高真正的坑在版本组合Spark 3.4 配 Scala 2.12spark-sql-kafka 的 artifactId 里必须带 _2.12装错会在作业启动时报 ClassNotFound。spark-submit \ --class com.example.RealtimeProfileJob \ --master yarn \ --deploy-mode cluster \ --num-executors 24 \ --executor-cores 4 \ --executor-memory 8g \ --conf spark.sql.shuffle.partitions96 \ --conf spark.sql.streaming.schemaInferencefalse \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.4.1 \ realtime-profile-0.1.0.jarexecutor-cores 固定 4 而不是 8单个 executor 内并发 task 太多时Kafka 拉取、GC、shuffle 会互相挤占 CPU4 核是流作业最常见的起点。spark.sql.shuffle.partitions 取 96等于 24 个 executor 乘 4 核窗口聚合的 shuffle 落到 96 个分区单批 50 万条以内每个分区数据量可控。配置项推荐值作用--executor-cores4单 executor 并发数过大会放大 GCspark.sql.shuffle.partitionsexecutor 总核数聚合与 join 的 shuffle 分区数maxOffsetsPerTrigger500000 起调单批最大消费量背压主开关3.2 读取 Kafka 并解析事件schema 先于代码定稿读 Kafka 的标准写法是把 value 反序列化再 from_json。schema 要在写代码前定稿每改一次字段checkpoint 里的状态和 schema 对不上最省事的处理是删 checkpoint 重来——这就是把 schema 定稿放在第一步的原因。import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ val spark SparkSession.builder() .appName(realtime-user-profile) .config(spark.sql.shuffle.partitions, 96) .getOrCreate() val eventSchema new StructType() .add(user_id, StringType, nullable false) .add(event_time, StringType) // 先按字符串读时区格式更可控 .add(event_type, StringType) .add(item_id, StringType) .add(scene, StringType) .add(props, MapType(StringType, StringType)) val parsed spark.readStream .format(kafka) .option(kafka.bootstrap.servers, kafka1:9092,kafka2:9092) .option(subscribe, page_views,orders,searches) .option(startingOffsets, earliest) .option(maxOffsetsPerTrigger, 800000) .option(failOnDataLoss, false) .load() .selectExpr(CAST(key AS STRING) AS user_key, CAST(value AS STRING) AS json) .select(from_json($json, eventSchema).as(e)) .select(e.*) .withColumn(event_time, to_timestamp($event_time, yyyy-MM-ddTHH:mm:ssXXX)) .filter($user_id.isNotNull $event_time.isNotNull)几个参数的意义startingOffsets 只在首次启动时生效之后以 checkpoint 记录为准所以新作业才设 earliestfailOnDataLoss 设成 falseKafka 里 offset 因过期被清理时作业不会直接崩而是从当前可用位点继续event_time 用显式格式转时间戳比放任 Spark 自动推断可靠带时区的 ISO8601 在 from_json 里直接按 TimestampType 解析经常失败。注意from_json 解析失败不会报错坏行会变成 null 字段所以必须 filter event_time 和 user_id。3.3 Spark DataFrame 聚合水位线、滑动窗口与分钟级标签生成聚合是标签计算的核心。按“用户 5 分钟窗口”把行为折叠成一组指标再用指标更新画像。对运营场景pv、下单数、去重商品数是三个最基础的指标。val profileAgg parsed .withWatermark(event_time, 10 minutes) .groupBy( $user_id, window($event_time, 5 minutes, 5 minutes) ) .agg( sum(when($event_type page_view, 1).otherwise(0)).as(pv), sum(when($event_type order, 1).otherwise(0)).as(order_cnt), countDistinct(when($event_type page_view, $item_id)).as(view_items), max($event_time).as(last_active) )watermark 的 10 分钟表示容忍事件最晚 10 分钟到达超过界限的旧数据直接丢弃它同时控制状态清理节奏窗口结束时间早于 watermark 的 key 会被移除这就是状态膨胀问题的源头。窗口长度和滑动步长都是 5 分钟输出是互不重叠的切片想平滑可以改成 window 10 分钟、slide 5 分钟代价是窗口状态量翻倍。countDistinct 在流式状态下按 key 保存整个 set浏览量大用户轻松吃掉几 MB 内存。画像允许少量误差换成 approx_count_distinct($item_id, 0.01)状态从 set 变成 HyperLogLog内存降一个量级误差控制在 1% 上下这是 Spark 内存调优里收益最大的一处改动。3.4 foreachBatch 双写 Redis 与 HBase连接复用与幂等聚合输出要写到多个下游。foreachBatch 把每批结果包装成普通 DataFrame在这批上做任意操作是多路写出的标准答案。profileAgg.writeStream .foreachBatch { (batchDF: org.apache.spark.sql.DataFrame, batchId: Long) batchDF.persist(org.apache.spark.storage.StorageLevel.MEMORY_AND_DISK) writeRedisIncremental(batchDF) // 在线标签profile:{user_id} writeHBaseHistory(batchDF) // 全量历史user_tags 宽表 publishUpdateEvents(batchDF) // 推送型标签变化 batchDF.unpersist() } .option(checkpointLocation, hdfs:///user/spark/checkpoint/realtime-profile) .trigger(org.apache.spark.sql.streaming.Trigger.ProcessingTime(60 seconds)) .start() .awaitTermination()写 Redis 不能在主函数里逐条建连接用 foreachPartition 让每个 task 复用一条连接pipeline 把一批 hset 合并成一次网络往返import redis.clients.jedis.Jedis def writeRedisIncremental(df: DataFrame): Unit { df.foreachPartition { rows val jedis new Jedis(redisHost, redisPort) rows.foreach { r val key profile: r.getAs[String](user_id) val pipe jedis.pipelined() pipe.hset(key, pv_5m, r.getAs[Long](pv).toString) pipe.hset(key, last_active, r.getAs[java.sql.Timestamp](last_active).toString) pipe.expire(key, 7 * 24 * 3600) pipe.sync() } jedis.close() } }ttl 设 7 天只保留近一周活跃用户冷用户自动淘汰。HBase 那路的关键是把 batchId 写进行里做幂等相同 batchId 的数据重放时不重复累加。4. 画像标签体系与存储设计让实时用户画像变得可查询4.1 标签三层拆分事实标签、规则标签、模型标签各归各位标签类型计算方式典型时效例子事实标签行为直接统计无加工秒到分钟5 分钟 pv、最近浏览商品、最后活跃时间规则标签阈值或条件判定分钟级高活跃、加购未下单、近 7 日复购率 0.5模型标签离线模型打分小时到天购买概率、流失概率、兴趣向量三层边界不划清楚什么标签都想塞进流里作业会被拖垮。事实标签是流的本职行为进来就算延迟越低越好。规则标签依赖连续窗口的累加状态“近 7 日复购率”要读历史订单不能只靠 5 分钟窗口常见做法是流里维护 7 日滚动计数或每天用批次预计算再合并。模型标签最不该实时化一个打分嵌进流里要引入特征服务依赖一次抖动就阻塞整条链路离线每天跑一次写进画像表查询时和实时标签合并返回。这套分层本质上是个标准的 spark数据分析案例同一份行为明细不同时效口径各取所需。4.2 存储选型对比Redis、HBase、ES 在画像链路里各管一段存储读模式写模式定位Redis按 user_id 直查毫秒级HSET / INCR / EXPIRE在线查询、实验取数HBaserowkey 扫描Put / Increment全量标签历史、回放重建ES / ClickHouse条件聚合检索批量写运营圈选、画像分析报表三者的边界可以简化为Redis 管在线HBase 管事实ES / ClickHouse 管检索。Redis 只放近一周活跃用户全量塞进去既不经济还会因内存上限被迫淘汰有用标签。HBase 的 rowkey 直接用 user_id列族 tags 下每列一个标签名值是 {value, ts} 的 JSON按用户读一次 get 拿全量按标签扫走全表这类低频查询丢给 ES 异步同步。ES 只当索引层查询命中后回 HBase 拉明细不要让 ES 当事实源它的更新延迟和丢失窗口在画像这种高频覆盖场景里藏不住。4.3 标签更新三模式覆盖写、原子累加与 TTL 淘汰三个模式对应三种语义混用是画像库最常见的坑。import org.apache.hadoop.hbase.TableName import org.apache.hadoop.hbase.client.{Connection, Increment} import org.apache.hadoop.hbase.util.Bytes def incrementTag(conn: Connection, userId: String, tag: String, delta: Long): Unit { val table conn.getTable(TableName.valueOf(user_tags)) val inc new Increment(Bytes.toBytes(userId)) inc.addColumn(Bytes.toBytes(cnt), Bytes.toBytes(tag), delta) table.increment(inc) table.close() }覆盖写用于事实标签比如 last_active直接 hset 新值覆盖。原子累加用于只增不减的计数器比如近 7 日订单数HBase 的 increment 是行级原子操作并发写不会互相覆盖但它只能保证累加不错乱不能保证不重复——同一个事件被流重放两次数字照样翻倍。所以进累加前要先做事件去重对 event_id 做布隆过滤或把 batchId 写进行里做幂等判断。TTL 用于短期标签session 级临时标签用 expire 自动清掉删除不要物理删行写入 tombstone 标记等下游读完再清理避免读到半新半旧的状态。4.4 补偿机制回放 Kafka 历史数据而不是重写一套批任务补偿逻辑不要做成独立的 spark etl脚本而是复用线上同一个 jar 的同一段聚合代码只把输入从 readStream 换成 read。Spark 3.3 起 Kafka 源支持按时间戳指定读取范围回放窗口精确到分钟。val replay spark.read.format(kafka) .option(kafka.bootstrap.servers, kafka1:9092,kafka2:9092) .option(subscribe, page_views) .option(startingOffsetsByTimestamp, {page_views:{0:1736755200000,1:1736755200000}}) .load()startingOffsetsByTimestamp 的值是分区号到 Unix 毫秒时间戳的映射实际主题有多少分区就得写多少3.3 之前的版本没有这个选项先用 kafka-get-offsets 工具按时间查各分区 offset再回填 startingOffsets。回放产出的是某时间段的标签增量与线上 Redis 对比时只看 batchId 大于当前值的部分覆盖更新。注意回放作业不要用线上同一个 checkpoint否则位点会被冲掉。5. 参数调优与准确性验证让实时画像健康地跑到上线之后5.1 背压控制与 Spark 内存配置maxOffsetsPerTrigger 之外还要看什么背压的第一道开关是 maxOffsetsPerTrigger限制一个 trigger 最多消费多少条消息。推荐从 50 万起调跑一个小时后看 Spark UI Streaming 页的 Input Rate 和 Kafka 消费组 laginput rate 不再波动、lag 归零逐步加量直到 CPU 使用率接近 75% 就停。调大 executor 数量经常没效果因为单批瓶颈在聚合和写出不在消费。内存侧executor-memory 8g 配 4 核是常用起点堆外留 512m 到 1g 给 Kafka consumer 和 HBase 客户端。GC 频繁时先看 State 大小而不是盲目加内存多数情况是 watermark 没生效导致状态只进不出。5.2 两个高频故障数据倾斜与状态膨胀怎么定位倾斜的典型表现是一个 task 跑不完日志里出现的全是同一批 user_id。头部用户行为量是普通用户的几十到上百倍5 分钟窗口里照样挤爆单分区。常见做法是加盐把 user_id 拼上 0 到 N 的随机后缀拆开聚合最后按真实 user_id 合并画像这类最终要精确到人的场景加盐只用于中间聚合输出前必须还原。另一个做法是把热 key 单独抽出来路由到独立作业避免拖慢全链路。状态膨胀看 Spark UI 查询详情页的 State 行规模持续上升、不随时间回落就是 watermark 或过期时间设置失效。5.3 回放对账与标签抽检验证实时画像准确率的两个手段准确率验证比性能调优更难做。第一个手段是回放对账把线上某小时 Kafka 全量数据用同一套聚合代码离线重放得到这一小时的标签增量与实时结果按 user_id 对齐比较 tag 值分布。差异超过 2 的用户数占总量比例低于 1% 才算合格这个 SQL 可以直接用在对比表上SELECT tag_name, count_if(real_time IS NOT NULL AND replay IS NOT NULL AND abs(real_time - replay) 2) AS matched, count(*) AS total FROM tag_compare GROUP BY tag_name;第二个手段是抽样标注每批随机抽 200 个 user_id由运营人工核对标签是否符合最近行为。整体分布一致但抽检命中率低几乎都是规则口径问题比如加购判定漏了某种入口。长期有效的技巧是把“标签生成时间 - 行为时间”的差值作为字段写进画像表每周统计 p50 / p95这个指标比任何监控面板更能说明实时性有没有退化。上线后盯 lag 用一条命令kafka-consumer-groups --bootstrap-server kafka1:9092 --group realtime-profile --describe连续几个周期 lag 不清零优先回调 maxOffsetsPerTrigger 而不是加 executor加 executor 只会放大状态和网络开销。本文还有配套的精品资源点击获取

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价