简介:这份源码资源面向大数据与电商方向的开发者、数据挖掘学习者,提供一套基于Spark的电商用户画像完整实现,可用于理解用户行为分析、标签体系构建与个性化推荐的技术落地路径。压缩包共462个文件,约13.45MB,以296个class编译文件、70个scala源文件、20个java源文件为核心,辅以properties、xml、json等配置与数据交换文件,以及js、css、html前端展示文件,另有jar包、字体与图标等资源,覆盖从数据处理到界面呈现的完整链路。目录按tags-model、tags-web、tags_ml、tags-etl等模块划分,分别对应模型定义、前端展示、机器学习与ETL数据清洗环节,结构清晰,便于按模块研读。目前已有339人学习下载,适合希望掌握Spark分布式计算、Scala函数式编程及用户画像建模的读者参考借鉴,也可作为课程设计或项目实战的对照方案。
1. 从一份电商用户画像源码说起:Spark 到底在算什么
打开任何一个电商后台,你都能看到「猜你喜欢」「人群包」「高价值客户」这些标签。它们背后往往就是一套用户画像系统。而「基于 Spark 的电商用户画像数据挖掘项目源码」这个标题,说的正是用 Spark 把订单、行为、商品等原始数据,加工成可查询、可投放的人群标签,并把这套流程写成可复现的工程代码。它解决的不是算法炫技问题,而是数据量大到单机跑不动时,怎么把标签算出来、存下来、用起来。适合有 SQL 基础、想从单机 pandas 转向分布式数据挖掘的工程师,也适合需要给运营团队交付人群包的数据开发。源码的价值在于把「口径、调度、存储、更新」四件事串成一条能跑的链路,而不是只给你一个模型文件。
2. 画像标签体系怎么拆:从原始表到人群包的字段映射
2.1 先定标签口径,再谈 Spark 代码
很多项目翻车的起点不是代码写错,而是口径没对齐。运营说「高价值客户」,可能指近 30 天消费金额前 10%,也可能指累计消费超过 5000 元。这两种口径在 Spark 里写出来的 SQL 完全不同。我一般会先拉一张标签定义表,字段包括:标签编码、标签名称、统计周期、计算逻辑、数据来源、更新频率。这张表不写进代码,但它是后面所有groupBy和agg的依据。
以电商场景为例,常见标签分四层:
| 层级 | 示例标签 | 主要数据源 | 更新频率 |
|---|---|---|---|
| 人口属性 | 性别、年龄段、城市等级 | 用户注册表 | 每天 |
| 消费能力 | 近30天GMV、客单价、消费频次 | 订单表 | 每天 |
| 行为偏好 | 近7天点击品类、收藏品类 | 行为埋点表 | 每小时 |
| 风险标签 | 退款率、投诉次数 | 售后表 | 每天 |
这张表决定了后面 Spark 任务的分区策略。人口属性数据量小,可以广播;行为数据量大,必须按用户 ID 分桶。源码里如果把这四层混在一个 Job 里跑,很容易出现数据倾斜,后面会讲怎么排查。
2.2 用 Spark SQL 把订单表聚合成消费标签
假设原始订单表ods_order有字段:user_id、order_id、pay_amount、pay_time、order_status。要算近 30 天消费金额和消费频次,最小可跑通的代码如下:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, sum, datediff, current_date spark = SparkSession.builder \ .appName("user_profile_consumption") \ .config("spark.sql.shuffle.partitions", "200") \ .enableHiveSupport() \ .getOrCreate() # 读取订单表,只取已支付订单 order_df = spark.table("ods_order") \ .filter(col("order_status") == "PAID") \ .filter(datediff(current_date(), col("pay_time")) <= 30) # 按用户聚合,同时算金额和频次 consumption_df = order_df.groupBy("user_id").agg( sum("pay_amount").alias("gmv_30d"), count("order_id").alias("order_cnt_30d") ) # 写入画像宽表,按 user_id 分区便于后续点查 consumption_df.write.mode("overwrite") \ .partitionBy("user_id") \ .saveAsTable("dws_user_consumption_30d")这段代码的逻辑很直白:先过滤有效订单和时间窗口,再按用户聚合。关键参数是spark.sql.shuffle.partitions,默认 200 在数据量小的时候会拖慢速度,数据量大时又不够。我一般会按数据量GB * 2估算,或者直接开 AQE(自适应查询执行)让 Spark 自己调。partitionBy("user_id")是为了后面按用户查标签时能走分区裁剪,但要注意如果用户 ID 基数太大,会产生大量小文件,需要配合repartition或写入后合并。
2.3 行为偏好标签:从埋点表到品类偏好向量
行为埋点表通常长这样:user_id、item_id、cate_id、behavior_type、event_time。要算「近 7 天点击最多的品类」,可以用窗口函数:
from pyspark.sql.window import Window from pyspark.sql.functions import row_number, desc behavior_df = spark.table("ods_behavior") \ .filter(col("behavior_type") == "CLICK") \ .filter(datediff(current_date(), col("event_time")) <= 7) # 按用户和品类聚合点击次数 cate_cnt = behavior_df.groupBy("user_id", "cate_id").agg( count("item_id").alias("click_cnt") ) # 取每个用户点击次数最多的前3个品类 window_spec = Window.partitionBy("user_id").orderBy(desc("click_cnt")) top_cate = cate_cnt.withColumn("rn", row_number().over(window_spec)) \ .filter(col("rn") <= 3) \ .groupBy("user_id") \ .agg(collect_list("cate_id").alias("top_cate_7d")) top_cate.write.mode("overwrite").saveAsTable("dws_user_cate_pref_7d")这里用row_number而不是rank,是因为并列时rank会跳号,导致取前 3 可能拿到 4 个品类。collect_list的顺序不保证稳定,如果下游对顺序敏感,需要额外加排序字段。行为数据往往比订单大一个量级,所以这个 Job 建议单独调度,不要和订单标签混在一起跑。
3. 源码工程怎么落地:目录结构、调度与存储选型
3.1 一份能跑的 Spark 画像源码通常长什么样
拿到「基于 Spark 的电商用户画像数据挖掘项目源码」后,先别急着跑main。我一般会先看目录结构,判断它是不是按分层做的。一个可维护的工程通常有这些目录:
src/ main/ scala/ 或 python/ etl/ # 原始数据清洗 label/ # 标签计算 profile/ # 画像宽表组装 utils/ # 工具类,如 SparkSession 封装 resources/ config/ # 标签口径配置、数据库连接 sql/ # 核心 SQL 脚本 test/ # 单元测试 pom.xml 或 requirements.txt如果源码把所有逻辑塞在一个Main.java或main.py里,维护成本会很高。标签口径一变,就要改代码重新打包。更好的做法是把标签逻辑写成配置或 SQL 文件,代码只负责读取配置、拼 SQL、调度执行。这样运营改口径时,数据开发只需要改配置,不用重新编译。
3.2 调度参数怎么设:内存、并行度和 Shuffle
Spark 任务跑得慢,八成是资源参数没设对。以下是我在画像项目里常用的配置,按数据量从大到小给三档:
| 数据量级 | executor 内存 | executor 核数 | shuffle partitions | 备注 |
|---|---|---|---|---|
| < 100GB | 4g | 2 | 100 | 开 AQE |
| 100GB-1TB | 8g | 4 | 400 | 开 AQE + 倾斜处理 |
| > 1TB | 16g | 4 | 800 | 分桶 + 倾斜处理 |
spark.executor.memory不是越大越好。超过 16g 后,GC 停顿会明显变长。如果任务频繁 Full GC,优先考虑加 executor 数量,而不是加单机内存。spark.sql.adaptive.enabled=true在 Spark 3.x 里基本是必开项,它能自动合并小分区、处理倾斜 Join。但 AQE 不是万能药,如果倾斜 key 的数据量超过单分区内存,还是会 OOM,这时候需要手动加盐。
3.3 画像结果存 Hive、HBase 还是 ClickHouse
标签算完后存哪里,取决于查询模式。运营圈人群包通常是「按标签组合筛选用户 ID」,这种场景用 ClickHouse 或 Doris 点查最快。但如果下游是推荐系统,需要按用户 ID 批量拉取标签向量,HBase 或 Redis 更合适。Hive 适合离线分析,但不适合高并发点查。
我一般会做两层存储:明细标签写 Hive 分区表,供 T+1 离线分析;同时把最新一天的画像宽表同步到 ClickHouse,供运营实时圈人。同步可以用 Spark 的 ClickHouse 连接器,或者写一个简单的 JDBC 写入任务。注意 ClickHouse 的ReplacingMergeTree引擎需要指定ver字段,否则重复写入不会自动去重。
4. 避坑与排查:画像任务跑不动时先看这 5 条
4.1 数据倾斜:某个用户 ID 占了 80% 的数据
现象:任务卡在最后一个 Stage,大部分 Task 已完成,少数几个 Task 跑了几小时还没结束。
原因:某个用户 ID 是爬虫或测试账号,行为数据量远超正常用户。groupBy("user_id")时,这个 key 落到一个分区,单点压力过大。
解决:先跑df.groupBy("user_id").count().orderBy(desc("count")).show(10)确认倾斜 key。如果是无效账号,直接在 ETL 层过滤;如果是正常大客户,用加盐法:给 key 拼接随机后缀,聚合后再去掉后缀二次聚合。
4.2 小文件过多:写入 Hive 后 NameNode 压力大
现象:partitionBy("user_id")写入后,HDFS 上出现几十万个小文件,后续查询启动 MapTask 极慢。
原因:用户 ID 基数大,每个分区至少一个文件,Spark 并行度高时每个 Task 写一个文件。
解决:写入前用repartition(200, "user_id")控制文件数,或者写入后跑一个合并任务。更推荐用bucketBy分桶表,但分桶数要提前定好,改起来麻烦。
4.3 内存溢出:collect_list 把数据全拉到 Driver
现象:任务报java.lang.OutOfMemoryError: GC overhead limit exceeded,Driver 日志显示collect相关调用。
原因:代码里用了collect_list或collect_set把每个用户的所有行为聚成一个数组,数据量大的用户直接把 Driver 撑爆。
解决:限制聚合数量,比如只取 Top 10;或者改用approx_count_distinct做近似统计。如果必须保留全量,写到 HDFS 而不是 collect 回 Driver。
4.4 时间窗口算错:时区导致 T+1 标签漏数据
现象:每天凌晨跑任务,发现前一天 23:00 后的订单没算进去。
原因:Spark 默认用 UTC 时区,而业务数据是北京时间。current_date()在 UTC 下比北京时间晚 8 小时。
解决:在 SparkSession 里设置spark.sql.session.timeZone=Asia/Shanghai,或者所有时间过滤都用明确的时间字符串,不依赖current_date()。
4.5 标签覆盖写导致历史数据丢失
现象:某天任务重跑后,前几天的标签数据被覆盖,运营反馈人群包人数变少。
原因:写入模式用了overwrite,且分区字段没包含日期,导致全表覆盖。
解决:画像表按dt分区,写入时用partitionBy("dt")+mode("overwrite"),这样只覆盖当天分区。如果标签需要保留历史快照,用insertInto追加,并加dt分区字段。
5. 进阶技巧:用 Spark 做标签回溯与人群包版本管理
标签算出来只是第一步,运营经常问「上个月的高价值客户和这个月差多少」。这就需要标签回溯能力。我的做法是在画像宽表里加两个字段:dt(数据日期)和version(标签版本号)。每次口径调整,version 加 1,旧版本数据保留。这样既能对比不同版本的人群差异,也能在口径出错时快速回滚。
具体实现上,可以用 Spark 的merge或insert overwrite按dt分区写入。查询时用where dt='2024-01-01' and version=2就能拿到指定版本的人群包。如果数据量太大,只保留最近 90 天的历史分区,更早的归档到对象存储。
另一个实用技巧是人群包导出。运营通常要的是用户 ID 列表,而不是宽表。可以用 Spark 把筛选后的user_id写成 CSV 或 Parquet,注意控制文件大小在 128MB 左右,方便下游导入。如果人群包要推送到广告平台,还需要做 ID 映射,比如把内部 user_id 转成设备 ID 或手机号哈希。这一步涉及数据安全,建议单独建一个映射表,不要直接在画像表里存明文手机号。
验证标签质量时,我习惯抽一批用户做人工核对。比如随机选 100 个「高价值客户」,看他们的订单金额是否真的排在前 10%。如果偏差超过 5%,就要检查口径或数据源。这个习惯帮我提前发现过好几次数据同步延迟的问题。希望帮到你。
本文还有配套的精品资源,点击获取