ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Spark实时用户画像系统实践:流式计算与特征存储全解析

Spark实时用户画像系统实践:流式计算与特征存储全解析 简介用户画像长期依赖离线T1批处理但实时推荐、在线风控和运营活动要求特征在分钟级甚至秒级生效传统的离线数仓模式已难以支撑这类低延迟场景。流式计算作为一种基于事件驱动、持续处理增量数据的计算范式天然契合实时特征生成的需求。在技术实现上通过Spark Structured Streaming消费Kafka中的用户行为日志利用状态存储维护用户历史偏好并结合事件时间窗口与水印机制即可构建一套实时标签计算管道。其技术价值在于相比离线批处理流式画像能及时捕捉用户最新意图有效提升推荐命中率与策略响应速度。该方案广泛适用于电商精准运营、个性化推荐、反欺诈风控等业务尤其适合已具备离线数仓基础、正从T1向实时数仓演进的团队。围绕Spark实时画像系统的工程落地可以从标签分层设计、ID归一化、流式管道搭建、特征存储选型到线上性能调优形成一套完整且可复用的实践路径。1. 实时用户画像离线T1的标签为什么必须搬到流上算用户画像系统听起来是离线数仓的活——T1把标签算好第二天早上查报表。但等你真做线上运营活动、实时推荐或风控策略时一套延迟一天的画像基本是废的用户半小时前刚搜过某类商品系统还在用三天前的标签决定给他推什么这不是画像是黑匣子。基于Spark的实时用户画像分析系统就是把这一套拉平到分钟级Kafka进实时行为Structured Streaming做流式计算状态存储保留历史偏好计算出的标签全部落库供在线服务查询。适合谁适合有Spark离线数仓基础、想把手头标签实时化的团队也适合做实时数仓项目但不知道怎么设计特征存储的开发者。这里要预先说清楚一个边界实时画像不是把离线任务改成流式语法就跑真正的成本在标签体系设计、状态管理和特征存储选型。接下来我按自己做过的一套方案从标签拆解讲到流式管道再讲到查询侧落地最后是踩坑记录和调优经验。2. 把画像拆成标签分层设计与ID归一化的三条经验很多团队上手就写聚合代码event进来直接count最后发现标签之间互相矛盾——按设备ID算一个画像按手机号又算出一个画像运营根本不知道该信哪个。问题不在Spark在标签体系设计阶段就埋了雷。2.1 标签分类事实标签、规则标签、模型标签的处理差别我在项目里把标签分成三类三类对实时性要求完全不同。事实标签是最简单的一类用户今天浏览了几个商品、加购几次、昨日活跃天数直接对行为事件做累计或计数即可实时性要求最高延迟控制在分钟级。规则标签需要写业务条件比如“近7天加购但未下单”“高活跃低转化用户”这类标签依赖多日数据光看当日流算不出来必须结合历史聚合结果。模型标签更复杂通常是离线训练好的模型对用户打分实时部分只负责把最新特征拼进模型输入比如推荐场景的CTR预估分数。三类的实时策略我建议分开事实标签走纯流式计算规则标签用“离线预聚合 流式增量更新”两条链路合并模型标签则落在特征服务层由在线推理服务按需拉取。很多方案把三类混在一个Pipeline里最后要么延迟压不下来要么标签口径对不上。2.2 ID归一化cookie、手机号、设备ID怎么合并成统一用户标识用户在看一个Spark实时画像项目时最容易忽略的却是最先要解决的ID映射。同一个用户可能用cookie访问Web端、用设备ID访问App端登录后才绑定手机号或用户ID。如果流里直接拿cookie当user_id画像会碎成几份后面所有统计都是错的。我一般会在接入层做一张“ID映射表”来归一化上游事件必须带一个统一的user_id才能进入后续计算。映射关系用HBase维护Web/App端行为先查映射拿到统一ID再写入Kafka。流式计算里只认统一ID不关心原始ID长什么样。另一种常见做法是把这段映射逻辑放在Structured Streaming的流里用状态存储维护减少一次在线查询但会显著增大状态规模我通常只在量级不大时才这么做。实时画像的数据质量七成靠这个ID归一化而不是靠SQL写得多花哨。3. 用Structured Streaming消费Kafka从JSON到画像特征的实时管道接下来进入核心实现Spark消费Kafka里的用户行为解析JSON做窗口聚合维护用户偏好状态最后把结果写出去。这套Pipeline就是搜“spark中读取json”“spark集群搭建”的从业者最想看到的落地样板。3.1 读Kafka的最小可用配置从JSON解析到schema强制Kafka里的行为日志一般是JSON字符串Spark读取时先拿到二进制value再通过from_json解析成结构化列。下面是一段可以直接跑起来的最小代码包含了流式消费的必配参数# streaming_user_profile.py from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, current_timestamp from pyspark.sql.types import ( StructType, StructField, StringType, LongType, DoubleType ) spark SparkSession.builder \ .appName(realtime-user-profile) \ .config(spark.sql.streaming.schemaInference, false) \ .getOrCreate() # 端到端至少一次模式下用Kafka自身的事务保证幂等写 schema StructType([ StructField(user_id, StringType(), True), StructField(item_id, StringType(), True), StructField(category_id, StringType(), True), StructField(action, StringType(), True), StructField(price, DoubleType(), True), StructField(ts, LongType(), True), ]) raw spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, kafka-1:9092,kafka-2:9092) \ .option(subscribe, user_behavior) \ .option(startingOffsets, earliest) \ .option(maxOffsetsPerTrigger, 20000) \ .option(failOnDataLoss, true) \ .load() parsed raw.select( from_json(col(value).cast(string), schema).alias(d) ).select(d.*) # 把事件时间转成TimestampType后续窗口和水印都依赖它 parsed parsed.withColumn( event_time, (col(ts) / 1000).cast(timestamp) ).withWatermark(event_time, 5 minutes)这段代码有几个关键点。maxOffsetsPerTrigger是每批最多消费多少条实时任务刚上线时最容易被忽视——默认情况下Spark会尽可能多地拉取数据批次间隔拉长延迟秒变分钟。我习惯先设20000到50000跑几天观察批次处理耗时再逐步调大。failOnDataLoss设为true可以防止Kafka topic因日志清理导致数据丢失而静默跳过排查问题时会少一个黑匣子。from_json有一个隐蔽的坑JSON里多出字段时解析后的结构体会带上这个字段但你如果只在schema里定义了部分列多出的部分会被丢弃。反过来schema里有而数据偶尔没有的字段解析出来是null下游聚合时要做好null过滤。更推荐的方式是通过spark.sql.jsonSchemaInference.enabled做一次采样推断但生产环境我还是坚持手写schema原因是JSON里如果混入脏类型推断出的schema可能每天都不一样流任务重启一次口径变一次这种翻车我遇到太多次了。3.2 用mapGroupsWithState维护用户最近偏好窗口聚合适合统计类标签但“用户最近一次浏览的品类”“用户近1小时点击最多的商品”这类偏好标签需要的是按键分组的状态更新。这类需求用groupBy配合agg是做不了的——中间状态只存在于当前微批内批间不保留。Structured Streaming提供了mapGroupsWithState可以按用户ID维护一个自定义状态对象每个事件到达时更新它。from pyspark.sql.functions import expr from pyspark.sql.types import StringType # 定义一个简单的偏好状态保存最近点击最多的类别 def update_state(user_id, events, state): if state is None: from pyspark.sql.types import StructType, StructField, StringType # 实际生产环境建议用自定义类维护counts字典 state {category_counts: {}} prev_counts state.get(category_counts, {}) for e in events: c e[category_id] if c: prev_counts[c] prev_counts.get(c, 0) 1 # 取当前计数最高的类别作为这个用户的实时偏好 top_category max(prev_counts.items(), keylambda x: x[1])[0] if prev_counts else None state[category_counts] prev_counts return (user_id, top_category), state # 注意groupByKey返回的是KeyValueGroupedDataset streamed parsed \ .select(user_id, category_id, action) \ .groupByKey(lambda r: r[user_id]) \ .mapGroupsWithState(update_state, timeoutConfProcessingTimeTimeout, outputModeUpdate)用mapGroupsWithState时最容易犯的错是把状态做得太重。我见过有人把用户一个月内点击的所有item_id列表都塞进state状态存储直接爆掉executor频繁GC。正确姿势是只保存“够算当前标签”的最小状态——统计类标签保存计数序列类标签只保存最近几十条能表达偏好即可。实时偏好做到分钟级延迟就够了不要试图在流里复刻离线全量逻辑。3.3 输出与truncate到底该写几份数据计算完成了输出比计算本身更容易翻车。常见做法是写两份一份是明细聚合后的画像宽表按天分区一份是最新状态快照供在线查询覆盖更新。宽表给离线分析用状态快照给实时荐购和运营后台用。# 写宽表用append模式按事件时间自动落分区 query1 aggregated.writeStream \ .format(parquet) \ .option(path, hdfs:///warehouse/user_profile/daily) \ .option(checkpointLocation, hdfs:///checkpoint/user_profile_agg) \ .partitionBy(event_date) \ .trigger(processingTime30 seconds) \ .outputMode(append) \ .start() # 写状态快照用update模式覆盖存储用HBase或Redis这里示例是console query2 streamed.writeStream \ .format(console) \ .option(truncate, false) \ .trigger(processingTime30 seconds) \ .outputMode(update) \ .start() query1.awaitTermination()宽表按事件时间分区需要保证上游ts单位一致毫秒和秒混用会让分区时间提前或滞后1小时这个我后面在避坑章细说。checkpointLocation是流任务的后悔药记录消费位点和状态进度删了它任务重启后可能重复消费或丢状态。生产环境绝对不要放本地临时目录要放在HDFS或S3上否则executor重启状态全丢。4. 特征落地与在线查询ClickHouse还是HBase我为什么先写宽表标签算出来了如果查询侧扛不住实时画像照样白搭。这一章写存储选型和查询侧设计回到搜索“实时特征服务”的读者关心的核心问题。4.1 在线查询对存储的几个硬要求点查快、更新快、别让Spark背锅实时画像的特征存储和离线数仓不一样离线分析可以扫全表在线服务必须点查用户打开App时系统要在几十毫秒内拿到这个用户的标签。Spark在这里的任务是高频写入不是承接高并发查询让Spark直接扛查询流量是最常见的架构错位。存储选型上我对比过两个方向。HBase擅长点查和更新主键直接是user_id适合规则标签和偏好标签ClickHouse更适合分析型查询按列存储压缩率高但单行更新要依赖ReplacingMergeTree或CollapsingMergeTree这类机制不适合高频覆盖更新。我最终选了“宽表 双存储”的方案在线查询走HBase列族按标签大类拆分离线分析走ClickHouse宽表查询标签分布和圈人用。如果你的标签以“圈人”为主而不是单用户点查可以反过来ClickHouse为主HBase只存最新偏好。4.2 宽表设计要点主键顺序、列族划分、TTL回收实时画像宽表不要按标签一个一个建表后期JOIN会把你逼疯。我一般按用户维度建一张大宽表每类标签一组字段字段命名统一前缀比如fact_开头是事实标签rule_开头是规则标签model_开头是模型分数。主键设计上HBase的rowkey不要直接拼user_iduser_id如果是纯数字且自增会造成热点。常见做法是加盐MD5(user_id).substring(0,2) _ user_id把数据分散到多个region代价是查询时要拼同样的盐前缀。列族不要贪多HBase官方建议不超过3个列族每个列族是一个写路径。我会拆成两个cf_base存基础属性性别、年龄段、注册时间cf_behavior存行为偏好近1小时浏览数、近7日加购数、最近偏好品类。TTL必须设置行为标签只保留30天基础属性保留180天否则状态过期后存储只增不减集群磁盘告警早晚找上门。写入链路我用Structured Streaming的foreachBatch每个微批批量写HBase不要逐行put。批量Put合并RPC单条Put会让region server忙不过来。控制batch大小在2000到5000条左右写失败可以重试对HBase的压力也小得多。5. 实时画像避坑记录迟到数据、状态膨胀与倾斜的排查实例这套系统跑起来之后我的经验一半来自功能开发另一半来自线上问题排查。下面几条是排查案例里最典型的每一条都真实影响过数据质量和任务稳定性。5.1 窗口计算每天都对不上你以为的乱序不是乱序现象T1对账时发现实时算出的“今日活跃用户数”比离线数仓少3%~8%每天偏差不稳定。查了Kafka延迟、批处理耗时都没问题。原因事件时间戳精度不一致上游部分日志打印ts是秒部分是毫秒我统一除以1000转时间戳秒的时间戳被当成毫秒除时间直接变成1970年。另有一部分是客户端本地时钟不准用户手机时间慢10分钟事件时间比服务端接收时间晚被水印丢掉。解决在接入层统一做时间戳单位校准同时用水印加宽到“接收时间 - 事件时间”差值覆盖网络延迟和客户端漂移。withWatermark设5分钟只是兜底生产环境我按P99消费延迟来设一般是消费延迟的两倍。还要在接入层加一层字段校验时间戳超过当前接收时间2小时的直接丢弃并计数告警。5.2 状态存储越跑越大默认的mapGroupsWithState没有近线清理现象任务跑了20天HDFS上的state数据从几十GB涨到几百GBcheckpoint恢复时间从10分钟变成40分钟偶尔还OOM。原因mapGroupsWithState定义的state如果永远不触发超时每个用户的状态就一直保留。用户量是累计值老用户不会消失状态只增不减。Spark虽然配置了TTL但那是针对HBase的对流式状态不会自动清。解决使用ProcessingTimeTimeout在状态里记录最后更新时间连续N天不活跃的用户状态置空或移除。具体在update_state函数里返回超时事件配合timeoutConf一起用。清理规则要分标签类型行为偏好7天没更新就可以清基础属性30天没更新才清。还要监控state大小超过阈值要告警而不是等它OOM。5.3 某个渠道用户特征全部为空JSON schema推断骗了你现象新接入一个渠道的行为日志category_id字段怎么解析都是null同一份JSON用离线Spark读就是正常的。新渠道的画像特征整片为空。原因schema是我从离线表结构里复制的离线处理时Spark SQL有容错JSON解析失败会置null不报错。但该渠道的日志里category_id这一层嵌套比Web端多了一层外包装from_json路径对不上解析时静默失败。解决新渠道接入时必须先抽样打印raw JSON人肉确认字段层级再定schema。from_json解析结果是null的行单独过滤出来写dead letter队列不要静默丢弃。这属于实时画像链路数据质量最容易被忽略的一环Kafka消息体格式变了一次你的标签就空一片不停留在这个检查项上就会长期被蒙在鼓里。5.4 热点用户把executor打满倾斜不是加资源能解决的现象某个大促头部主播的粉丝用户行为量是普通用户的100倍这几个用户的数据全部hash到同一个executor该executor CPU跑满其他executor空闲整个流任务背压。原因groupByKey(user_id)对极端热点用户来说所有事件都到一个分区。加executor也没用数据热点不会因为机器多而分散。解决对超高活跃用户单独分桶。user_id加盐后再分组热点用户拆成多个子key分别计数最后再合并。或者给用户分VIP等级VIP用户的事件单独一个topic单独的流任务处理主链路只处理普通用户。这个方法麻烦但一百个用户吃掉十万级QPS的场景必须做。6. 把TP99打下来批量攒写、状态清理与执行内存的调优技巧实时画像上线后的路才走了一半性能调优是最后一步。分享几个我压测和线上调优时实际用的技巧。触发间隔不要盲目调到秒级。Spark Structured Streaming本质是微批processingTime设5秒和30秒对业务方感知差别不大但集群开销差5倍以上。我通常从30秒起跑持续观察batch处理耗时处理耗时低于触发间隔的70%再尝试降到15秒直到处理耗时的增长跟不上间隔缩短的趋势。批次大小用maxOffsetsPerTrigger控制的同时还要关注执行内存。流式任务spark.executor.memory的默认值通常不够我一般给到4GB以上因为状态存储和微批执行都在堆内。配合spark.memory.offHeap.enabledtrue和spark.memory.offHeap.size2g把部分状态挪到堆外GC压力会明显下降。堆外内存设置要留足余量否则直接抛OOM。写入存储侧的攒批也很重要。写给HBase/Redis这类在线存储时在foreachBatch里对同一批数据先按主键做一次去重再批量提交。Kafka内部重试可能产生重复记录流任务做了一次幂等处理后HBase的写放大能降30%以上TP99肉眼可见往下掉。验证方法上我自己习惯在压测环境造一份“正常用户 极端热点用户”混合数据分别检测吞吐、P95延迟、状态增长三项指标。跑通后再放灰度流量。这套系统从上线到现在我最深的教训是实时画像的技术难点从来不只在流式计算本身而在于把标签设计、ID归一化、状态清理和存储写入当成一个整体来做。最后一条经验是给所有想快速落地的人——先砍掉一半标签用最少的数据跑通全链路再逐步加标签否则你会在排错上消耗掉比开发还多的时间。希望帮到你。本文还有配套的精品资源点击获取
返回列表