
简介这份资源是面向大数据与电商方向学习者、算法工程师的Spark用户画像数据挖掘项目源码适合作为课程设计、毕业设计或技术进阶的实战参考。项目以Scala与Java为主力语言结合Spark分布式计算框架围绕用户行为与偏好构建精准画像支撑个性化推荐与精准营销场景。压缩包共462个文件约13.45MB包含296个class、70个scala、20个java等核心代码文件以及properties、xml、json等配置与数据交换文件另有jar依赖包和js、css、html等前端展示资源目录按tags-model、tags-web、tags_ml、tags-etl等模块划分清晰体现从ETL到模型训练再到Web呈现的完整链路。目前已有339人学习下载。读者可从中获取一套可运行的用户画像挖掘工程结构理解RFM、USG等标签模型的实现思路掌握HBase关系映射与ML工具类的组织方式并借鉴多模块协作与前后端配合的工程实践。1. 基于Spark的电商用户画像数据挖掘从标签体系到源码落地电商平台每天产生千万级行为日志用户点开一个商品、停留几秒、加购又放弃这些碎片化动作背后藏着可被量化的消费意图。基于Spark的电商用户画像数据挖掘项目源码核心要解决的就是把这些散落在Hive、Kafka、MySQL里的原始行为数据通过分布式计算打成一张张用户标签最终服务于推荐、营销和风控。我见过不少团队用Python单机跑用户画像数据量一过千万行就崩而Spark的分布式内存计算恰好能扛住这个量级。这个方向适合有基础SQL和Python能力、想从单机数据分析转向大数据工程的同学也适合需要给推荐系统补上画像层的中高级工程师。源码不是拿来直接跑的玩具而是理解标签计算链路的脚手架。2. 用户画像标签体系与Spark技术选型为什么不是Pandas2.1 电商画像的四类标签与计算粒度用户画像在电商场景里通常拆成四层统计类标签、规则类标签、挖掘类标签和实时类标签。统计类标签包括近7天点击次数、近30天客单价、累计下单金额这类标签用Spark SQL的窗口函数就能算计算粒度是天级或小时级。规则类标签比如“高价值用户”定义为近90天消费金额大于5000且下单次数大于10这类标签依赖业务规则用DataFrame的when/otherwise组合条件即可。挖掘类标签涉及RFM模型、用户生命周期阶段、商品偏好向量需要Spark MLlib做聚类或分类。实时类标签则要接Spark Structured Streaming从Kafka消费点击流做分钟级更新。计算粒度决定了存储方案。天级标签写入Hive宽表按dt分区每个用户一行列是标签名。小时级标签写入HBase或Redis供线上服务实时查询。我一般会把标签分成基础标签层和衍生标签层基础层只做清洗和简单聚合衍生层在基础层上做交叉计算这样回溯问题时能快速定位是原始数据脏了还是计算逻辑错了。2.2 Spark相比单机Pandas的硬性优势很多人问用户画像用Pandas加定时脚本不行吗数据量在百万行以内确实可以但电商场景的日活用户行为日志动辄上亿条Pandas单机内存根本装不下。Spark的核心优势在于三点第一数据分片后并行计算100个Executor能同时处理不同分区的数据第二Spark SQL的Catalyst优化器会自动做谓词下推和列裁剪读取Parquet文件时只加载需要的列第三Spark MLlib提供了分布式实现的KMeans、ALS、FPGrowth不用自己写分布式算法。选型时还要注意Spark版本。Spark 3.x对Structured Streaming的支持更完善AQE自适应查询执行能自动处理数据倾斜。如果团队还在用Spark 2.x建议至少升到2.4.8否则窗口函数和广播变量的性能差距很明显。存储格式优先选Parquet列式存储对画像宽表的聚合查询友好压缩比也比Text高3到5倍。2.3 最小可跑通的画像计算链路一个完整的画像计算链路包含数据接入层Kafka/Hive/MySQL、清洗层去重、过滤、字段标准化、标签计算层Spark SQL/MLlib、存储层Hive/HBase/Redis、服务层API查询。下面这段PySpark代码演示从Hive读取用户行为日志计算近7天点击次数和近30天下单金额两个基础标签。from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, sum, datediff, current_date, when # 初始化SparkSession开启AQE和动态分区 spark SparkSession.builder \ .appName(EcommerceUserProfile) \ .config(spark.sql.adaptive.enabled, true) \ .config(spark.sql.adaptive.coalescePartitions.enabled, true) \ .config(spark.sql.sources.partitionOverwriteMode, dynamic) \ .enableHiveSupport() \ .getOrCreate() # 读取行为日志表按dt分区过滤近30天数据 behavior spark.sql( SELECT user_id, item_id, behavior_type, dt FROM dwd.user_behavior_log WHERE dt date_sub(current_date(), 30) ) # 计算近7天点击次数 click_7d behavior.filter( (col(behavior_type) click) (datediff(current_date(), col(dt)) 7) ).groupBy(user_id).agg(count(*).alias(click_cnt_7d)) # 计算近30天不同下单天数去重 order_30d behavior.filter( col(behavior_type) order ).groupBy(user_id).agg( count(*).alias(order_cnt_30d), sum(when(col(behavior_type) order, 1).otherwise(0)).alias(order_flag) ) # 关联两个标签写入Hive宽表 profile click_7d.join(order_30d, onuser_id, howfull_outer).fillna(0) profile.write.mode(overwrite).insertInto(dws.user_profile_base)这段代码的逻辑是先从Hive分区表读取近30天行为日志利用分区裁剪减少扫描量然后分别计算点击和下单两个标签最后用full_outer join合并避免丢失只有点击没有下单的用户。参数方面spark.sql.adaptive.enabled开启后Spark会根据shuffle数据量自动调整分区数解决小文件问题partitionOverwriteMode设为dynamic写入Hive时只覆盖涉及的分区不会全表重写。注意fillna(0)要放在join之后否则可能把null值误填到不该填的列。3. 从原始日志到画像宽表清洗、聚合与RFM打分3.1 行为日志清洗的四个关键动作原始行为日志通常包含脏数据user_id为null、behavior_type出现未知枚举值、时间戳格式不统一、重复上报。清洗阶段要做四件事。第一过滤user_id为null或空字符串的记录这些数据无法归属到任何用户。第二校验behavior_type是否在允许集合内click、cart、fav、order、pay不在集合内的记录打标后写入脏数据表不要直接丢弃方便排查上游埋点问题。第三统一时间格式把字符串时间转成timestamp类型后续窗口函数才能正确计算。第四按user_id和item_id加时间戳去重防止客户端重复上报导致计数虚高。清洗后的数据写入dwd层按dt分区存储为Parquet。这里有个血泪经验不要用dropDuplicates()全字段去重数据量大时shuffle开销极高。正确做法是用row_number()窗口函数按user_id、item_id、behavior_type、timestamp分组取rn1的记录这样能精确控制去重粒度。3.2 RFM模型在Spark中的分布式实现RFM是用户画像里最经典的挖掘模型R代表最近一次消费时间间隔F代表消费频率M代表消费金额。单机实现很简单但分布式环境下要处理两个问题一是时间基准要统一所有Executor用同一个current_date不能各自取本地时间二是打分阈值要提前算好不能在每个分区里单独算分位数。我一般分两步走。第一步用Spark SQL算出每个用户的R、F、M原始值写入临时表。第二步用approxQuantile计算全局分位数再根据分位数给每个用户打1到5分。approxQuantile是近似算法误差可控在1%以内比精确分位数快一个数量级。# 计算RFM原始值 rfm_raw spark.sql( SELECT user_id, datediff(current_date(), max(order_date)) AS recency, count(DISTINCT order_id) AS frequency, sum(order_amount) AS monetary FROM dwd.order_detail WHERE dt date_sub(current_date(), 90) GROUP BY user_id ) # 用approxQuantile计算全局分位数阈值 quantiles rfm_raw.approxQuantile( [recency, frequency, monetary], [0.2, 0.4, 0.6, 0.8], 0.01 # 相对误差 ) # 根据阈值打1-5分recency越小分越高 def score_r(value): if value quantiles[0][0]: return 5 elif value quantiles[0][1]: return 4 elif value quantiles[0][2]: return 3 elif value quantiles[0][3]: return 2 else: return 1 # 注册UDF并应用 from pyspark.sql.functions import udf from pyspark.sql.types import IntegerType score_r_udf udf(score_r, IntegerType()) rfm_scored rfm_raw.withColumn(r_score, score_r_udf(col(recency)))这段代码的关键在于approxQuantile的第三个参数0.01表示相对误差为1%值越小越精确但计算越慢。实际项目中如果用户量在千万级这个参数设0.01足够如果上亿可以放宽到0.05。打分逻辑用UDF实现注意UDF内部不要引用外部变量比如quantiles否则序列化会出问题。正确做法是把阈值广播出去或者用when/otherwise链式表达。3.3 画像宽表的存储与更新策略画像宽表通常有几百列每天全量重写成本太高。我一般用两种策略增量更新和分区覆盖。增量更新适合标签只增不减的场景比如累计下单金额用merge into语法按user_id更新。分区覆盖适合天级快照标签比如近7天点击次数每天重新计算后覆盖当天分区。Hive表建议用ORC格式比Parquet的压缩比更高但Parquet对Spark的兼容性更好。如果画像表要供Impala或Presto查询选Parquet如果只供Spark和Hive查ORC更省空间。分区字段用dt桶字段用user_id分桶数设为Executor数的2到3倍这样join时能触发bucket map join避免shuffle。4. 避坑与排查Spark画像任务翻车的五个典型场景4.1 数据倾斜导致任务卡在99%现象Spark任务跑到最后几个task时卡住不动日志显示某个task处理的数据量是其他task的几十倍。原因user_id为null或某个热门商品被大量用户点击导致groupBy时某个key的数据量畸高。解决先过滤null key再对热门key加随机前缀打散计算完再去前缀聚合。或者开启AQE的倾斜处理设置spark.sql.adaptive.skewJoin.enabledtrue让Spark自动拆分倾斜分区。4.2 OOM频发但内存明明够现象Executor频繁OOM但监控显示内存使用率不到70%。原因可能是spark.sql.shuffle.partitions设得太大每个task数据量小但task数量多导致调度开销和内存碎片。也可能是UDF里创建了大对象没释放。解决把shuffle partitions从默认200调到Executor数乘以2到3倍UDF里避免创建大字典或列表用广播变量替代。4.3 标签计算结果每天波动异常现象同一个用户昨天的“近7天点击次数”是50今天变成5但用户行为没突变。原因时间窗口计算时用了current_date()而Spark任务跨天运行时不同Executor取到的时间可能不一致。解决在Driver端取一次时间作为常量广播到所有Executor或者直接用数据里的dt字段做基准不依赖系统时间。4.4 写入Hive后小文件泛滥现象画像表目录下有几万个几KB的小文件NameNode压力大查询也慢。原因每个Spark task写一个文件task数太多且没做合并。解决写入前用repartition或coalesce控制分区数或者在Hive端开启动态分区合并设置hive.merge.mapfilestrue和hive.merge.size.per.task。4.5 UDF序列化报错但代码没问题现象任务提交后报Task not serializable检查代码没发现闭包引用外部变量。原因UDF里引用了Driver端的对象比如SparkSession或数据库连接。解决UDF内部只使用传入的参数需要外部配置时用广播变量。另外Python UDF性能比Scala UDF差很多能用Spark SQL内置函数就别用UDF。5. 画像标签的线上验证与增量迭代技巧画像算完不是终点标签准不准、更新及不及时才是。我一般用三个手段做验证。第一抽样对比从Hive里随机抽1000个用户用Pandas单机复算一遍和Spark结果比对差异超过5%就要查逻辑。第二分布监控每天统计各标签的均值、中位数、空值率画趋势图突变时告警。第三A/B测试把画像标签接入推荐系统看点击率是否有提升这是最直接的业务验证。增量迭代时我习惯把标签计算拆成独立的任务每个任务只负责一类标签用Airflow或DolphinScheduler编排依赖。这样改一个标签的逻辑不会影响其他标签。另外标签版本要管理起来每次变更记录版本号和变更原因方便回溯。# 标签质量校验空值率和分布对比 def validate_profile(spark, table_name, key_col): df spark.table(table_name) total df.count() null_cnt df.filter(col(key_col).isNull()).count() null_rate null_cnt / total # 空值率超过10%告警 if null_rate 0.1: print(f警告{table_name} 的 {key_col} 空值率 {null_rate:.2%}) # 输出各标签的均值 numeric_cols [f.name for f in df.schema.fields if f.dataType.typeName() in (integer, double, long)] df.select(numeric_cols).describe().show()这段校验代码放在每天任务结束后跑空值率和均值写入监控表。注意describe()会触发全表扫描大表上要加采样比如df.sample(0.01)。最后说个我踩过的坑有次为了赶进度把画像标签直接写进Redis供线上查结果Redis内存爆了。后来改成HBase存全量、Redis只存热点用户才稳住。画像项目不是算完就完事存储选型和查询模式要提前想清楚。希望帮到你。本文还有配套的精品资源点击获取