简介:这是一套面向计算机相关专业毕业生、课程设计参与者及大数据学习者的Spark电商用户行为分析系统,包含完整源码与技术文档,可直接用于毕业设计、课程作业或项目实践。系统基于Spark分布式计算架构,支持实时分析与离线计算两种模式,核心功能覆盖用户点击流分析、购买行为模式识别、用户画像构建,并整合Spark MLlib协同过滤算法实现个性化推荐,借助Spark Streaming处理实时数据流,前端采用ECharts展示分析结果。资源包共286个文件,约1.28MB,以185个xml配置、40个java源码、47个zbak备份文件为主,另含properties配置、png图表及md说明文档,目录结构清晰,便于按模块检索学习。已有56人学习下载。项目在导师指导下完成并获学术评审99分,代码经严格验证可稳定运行,适合不同技术基础的用户参考部署与二次开发。
1. 从一份能跑通的 Spark 电商用户行为分析系统说起
电商后台每天沉淀的点击、加购、下单、支付日志,单机 Pandas 早就扛不住了,动辄几十 GB 的埋点文件让内存直接爆掉。Spark 电商用户行为分析系统源码与完整文档,讲的正是把这套分析链路从零搭起来:用 Spark 做数据清洗、会话切分、漏斗转化、RFM 分层,最后落到可复现的源码和一份能照着走的文档。它解决的不是"Spark 是什么",而是"我拿到一份电商行为日志,怎么在集群上跑出业务方要的指标"。适合两类人:一是要交课程设计或做大数据项目、需要一套完整可运行代码的在校同学;二是刚接手用户行为分析、想快速搭出可用管道的初中级数据工程师。下面按"环境怎么搭 → 数据怎么洗 → 指标怎么算 → 坑在哪 → 怎么验证"的顺序讲透。
2. 环境与数据准备:把 Spark 集群和电商行为日志先跑起来
2.1 本地伪分布式与集群模式的选型理由
做电商用户行为分析,第一步不是写业务代码,而是决定跑在哪。常见做法有三种:本地local[*]、Standalone 伪分布式、以及 YARN 上的集群模式。选型不看"哪个高级",看数据量和调试成本。
本地模式适合开发阶段,数据量在几 GB 以内,local[4]就能把逻辑跑通,改一行代码重启只要几秒。缺点是它不模拟真实的分区调度,很多在集群上才暴露的问题(数据倾斜、Executor 内存溢出)本地根本复现不出来。我一般会先用本地模式把清洗和指标逻辑写对,再切到集群验证。
Standalone 伪分布式是在一台机器上起 Master 和 Worker,能验证资源调度和并行度,适合课程设计这种"要展示集群能力但只有一台机器"的场景。YARN 模式才是生产常态,资源由 YARN 统一分配,--num-executors、--executor-memory这些参数才真正起作用。
提示:课程设计或演示场景,Standalone 伪分布式足够;真要对标生产,直接上 YARN 或 K8s 模式,别在伪分布式上纠结太久。
环境依赖上,JDK 8 或 11、Scala 2.12、Hadoop 3.x 是当前最稳的组合。Spark 3.x 对 JDK 17 的支持在部分组件上仍有兼容问题,新手别一上来就追最新 JDK。
2.2 用脚本把 Spark 环境拉起来的最小步骤
下面这段是本地和伪分布式通用的环境准备脚本,逻辑是下载解压、配置环境变量、改核心配置文件。参数按注释改。
#!/bin/bash # 下载并解压 Spark(版本按需替换,此处以 3.5.x 为例) wget https://archive.apache.org/dist/spark/spark-3.5.1/spark-3.5.1-bin-hadoop3.tgz tar -zxvf spark-3.5.1-bin-hadoop3.tgz -C /opt/ mv /opt/spark-3.5.1-bin-hadoop3 /opt/spark # 配置环境变量,写入 ~/.bashrc echo 'export SPARK_HOME=/opt/spark' >> ~/.bashrc echo 'export PATH=$PATH:$SPARK_HOME/bin:$SPARK_HOME/sbin' >> ~/.bashrc source ~/.bashrc # 伪分布式:配置 Master 和 Worker cd $SPARK_HOME/conf cp spark-env.sh.template spark-env.sh # 指定 Master 主机和端口,JAVA_HOME 按实际路径改 echo 'export SPARK_MASTER_HOST=localhost' >> spark-env.sh echo 'export SPARK_MASTER_PORT=7077' >> spark-env.sh echo 'export JAVA_HOME=/usr/lib/jvm/java-11-openjdk' >> spark-env.sh # 启动集群 $SPARK_HOME/sbin/start-master.sh $SPARK_HOME/sbin/start-worker.sh spark://localhost:7077逻辑说明:先解压到固定目录,避免路径漂移;环境变量写进~/.bashrc保证新终端可用;spark-env.sh里SPARK_MASTER_HOST决定 Master 绑定地址,SPARK_MASTER_PORT默认 7077,被占用就改。启动后用jps应能看到Master和Worker两个进程,Web UI 在 8080 端口。
参数说明:SPARK_WORKER_CORES控制单 Worker 可用核数,SPARK_WORKER_MEMORY控制内存,伪分布式下按机器实际资源给,别把全部内存分出去,否则系统本身会卡。
2.3 电商行为日志的字段结构与读入方式
电商用户行为日志通常是 JSON 行格式,一行一条事件。典型字段包括:user_id、item_id、category_id、behavior_type(pv/cart/fav/buy)、timestamp。读入时用 Spark SQL 的from_json配合显式 schema,比inferSchema快且稳。
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, LongType spark = SparkSession.builder \ .appName("ecommerce_user_behavior") \ .master("local[4]") \ .config("spark.sql.shuffle.partitions", "8") \ .getOrCreate() # 显式定义 schema,避免 inferSchema 全表扫描 schema = StructType([ StructField("user_id", StringType(), True), StructField("item_id", StringType(), True), StructField("category_id", StringType(), True), StructField("behavior_type", StringType(), True), StructField("timestamp", LongType(), True), ]) df = spark.read.schema(schema).json("hdfs:///data/user_behavior/*.json") df.createOrReplaceTempView("raw_behavior") df.printSchema()逻辑说明:master("local[4]")用 4 个核跑本地;spark.sql.shuffle.partitions默认 200,本地小数据量下会拖慢,改成 8 更合适。显式 schema 的好处是字段类型确定,后续timestamp转日期不会因为类型推断错误翻车。
参数说明:spark.sql.shuffle.partitions是血泪经验里最常被忽略的参数,集群上按数据量调,一般设成核数的 2 到 3 倍;本地调试设小。spark.default.parallelism影响 RDD 操作的默认并行度,和 shuffle partitions 是两套东西,别混。
3. 数据清洗与会话切分:把原始埋点变成可用行为流
3.1 脏数据识别与清洗规则设计
原始埋点里最常见的脏数据有四类:字段缺失(user_id为空)、时间戳异常(未来时间或 1970 年)、行为类型不在枚举内、重复上报。清洗不是无脑dropna,要按业务规则来。
from pyspark.sql import functions as F # 时间戳转标准时间,过滤异常时间范围 clean_df = df \ .filter(F.col("user_id").isNotNull()) \ .filter(F.col("behavior_type").isin("pv", "cart", "fav", "buy")) \ .withColumn("event_time", F.to_timestamp(F.from_unixtime("timestamp"))) \ .filter(F.col("event_time") >= F.lit("2024-01-01")) \ .filter(F.col("event_time") <= F.current_timestamp()) \ .dropDuplicates(["user_id", "item_id", "behavior_type", "timestamp"]) clean_df.createOrReplaceTempView("clean_behavior") print("清洗后条数:", clean_df.count())逻辑说明:先过滤空user_id,再限定行为枚举,from_unixtime把秒级时间戳转成可读时间,to_timestamp转成 TimestampType 便于后续按天聚合。时间范围过滤挡掉未来时间和远古脏数据。dropDuplicates按四个字段去重,防止同一事件重复上报导致漏斗虚高。
参数说明:dropDuplicates的字段选择很关键,选少了去重不干净,选多了会误删真实重复行为(比如用户短时间内多次浏览同一商品)。我一般按"用户+商品+行为+秒级时间"去重,秒级以内视为重复上报。
3.2 会话切分:30 分钟不活跃就断会话
用户行为分析的核心单位是"会话"(session),不是单条事件。会话切分规则通常是:同一用户相邻两次行为间隔超过 30 分钟,就切一个新会话。这个逻辑用窗口函数实现。
from pyspark.sql import Window # 计算同一用户相邻事件的时间差 w = Window.partitionBy("user_id").orderBy("timestamp") session_df = clean_df \ .withColumn("prev_ts", F.lag("timestamp").over(w)) \ .withColumn("gap", F.col("timestamp") - F.col("prev_ts")) \ .withColumn("is_new_session", F.when((F.col("gap").isNull()) | (F.col("gap") > 1800), 1).otherwise(0)) \ .withColumn("session_id", F.sum("is_new_session").over(w.rowsBetween(Window.unboundedPreceding, Window.currentRow))) session_df.createOrReplaceTempView("session_behavior")逻辑说明:lag取同一用户上一条事件时间,gap是秒级差值,超过 1800 秒(30 分钟)标记为新会话起点。sum累加is_new_session得到会话编号,同一会话内编号相同。这是会话切分的标准做法,比自写 UDF 快得多。
参数说明:1800 秒是电商场景的常见阈值,内容类产品可能用 15 分钟,直播场景可能用 5 分钟。阈值直接影响会话数和人均会话时长,业务方对不上数时先查这个值。
3.3 用 Spark SQL 做日期加减与时间维度扩展
行为分析离不开时间维度:按天、按周、按小时。Spark SQL 的日期函数能直接算,不用回 Python。
-- 按天、按小时聚合,并计算次日留存所需的日期偏移 SELECT user_id, session_id, DATE(event_time) AS dt, HOUR(event_time) AS hr, DATE_ADD(DATE(event_time), 1) AS next_dt, behavior_type FROM session_behavior WHERE event_time IS NOT NULL逻辑说明:DATE取日期,HOUR取小时,DATE_ADD(..., 1)得到次日日期,用于留存计算时把"当天行为"和"次日行为"关联。Spark SQL 的日期加减用DATE_ADD/DATE_SUB,和 MySQL 语法接近,迁移成本低。
参数说明:DATE_ADD第二个参数是天数,负数即减。跨月跨年由函数自动处理,不用手动判断。注意event_time必须是 TimestampType,字符串类型会报错。
4. 核心指标计算:漏斗、留存与 RFM 分层
4.1 转化漏斗:从浏览到支付的四步拆解
电商漏斗通常是"浏览 → 加购 → 收藏 → 支付"。用会话级或用户级聚合都行,关键是口径统一。
funnel = spark.sql(""" SELECT COUNT(DISTINCT CASE WHEN behavior_type = 'pv' THEN user_id END) AS pv_users, COUNT(DISTINCT CASE WHEN behavior_type = 'cart' THEN user_id END) AS cart_users, COUNT(DISTINCT CASE WHEN behavior_type = 'fav' THEN user_id END) AS fav_users, COUNT(DISTINCT CASE WHEN behavior_type = 'buy' THEN user_id END) AS buy_users FROM clean_behavior """) funnel.show()逻辑说明:用COUNT(DISTINCT user_id)按行为类型分别统计去重用户数,得到各环节人数。漏斗转化率 = 下一环节人数 / 上一环节人数。注意这里统计的是"有过该行为的用户",不是"行为次数",口径不同结论差很多。
参数说明:如果要算"严格漏斗"(必须按顺序发生),需要自关联或窗口函数判断行为先后,复杂度高但更准。多数业务方接受宽松口径,先确认再动手。
4.2 留存计算:次日、7 日、30 日留存
留存的核心是"某天新增用户,在之后第 N 天是否还有行为"。用自关联实现。
retention = spark.sql(""" WITH first_visit AS ( SELECT user_id, MIN(DATE(event_time)) AS first_dt FROM clean_behavior GROUP BY user_id ), daily_active AS ( SELECT DISTINCT user_id, DATE(event_time) AS active_dt FROM clean_behavior ) SELECT f.first_dt, COUNT(DISTINCT f.user_id) AS new_users, COUNT(DISTINCT CASE WHEN DATEDIFF(d.active_dt, f.first_dt) = 1 THEN d.user_id END) AS day1, COUNT(DISTINCT CASE WHEN DATEDIFF(d.active_dt, f.first_dt) = 7 THEN d.user_id END) AS day7, COUNT(DISTINCT CASE WHEN DATEDIFF(d.active_dt, f.first_dt) = 30 THEN d.user_id END) AS day30 FROM first_visit f JOIN daily_active d ON f.user_id = d.user_id GROUP BY f.first_dt """)逻辑说明:first_visit算每个用户首次活跃日期,daily_active算每日活跃用户,关联后用DATEDIFF判断间隔天数。DATEDIFF返回天数差,等于 1 即次日留存。
参数说明:留存口径要明确是"自然日"还是"24 小时",两者结果不同。DATEDIFF按自然日算,跨零点即算一天。数据量大时这个自关联会 shuffle,注意分区数。
4.3 RFM 分层:用窗口函数给用户打标签
RFM 是 Recency(最近一次消费)、Frequency(消费频次)、Monetary(消费金额)。电商行为数据里如果没有金额,可用购买次数近似。
rfm = spark.sql(""" WITH user_rfm AS ( SELECT user_id, DATEDIFF(CURRENT_DATE(), MAX(DATE(event_time))) AS recency, COUNT(CASE WHEN behavior_type = 'buy' THEN 1 END) AS frequency FROM clean_behavior GROUP BY user_id ) SELECT user_id, recency, frequency, NTILE(5) OVER (ORDER BY recency ASC) AS r_score, NTILE(5) OVER (ORDER BY frequency DESC) AS f_score FROM user_rfm WHERE frequency > 0 """)逻辑说明:recency越小越好,frequency越大越好,用NTILE(5)分五档打分。NTILE是等频分桶,保证每档人数接近,比固定阈值更稳。
参数说明:NTILE的桶数按业务需要定,5 档最常见。分档后可以组合成"高价值""流失预警"等标签,规则由运营定,代码只负责打分。
5. 避坑与排查:Spark 电商行为分析里最容易翻车的五件事
5.1 数据倾斜导致个别 Task 卡死
现象:任务跑到 99% 不动,Web UI 上某个 Task 处理的数据量远超其他。原因:某个热门商品或异常用户的行为记录特别多,groupBy或join时全压到一个分区。解决:先spark.sql.adaptive.enabled=true开自适应执行,让 Spark 自动处理倾斜;仍不行就对热点 key 加随机前缀打散,聚合后再合并。
5.2 Executor 内存溢出 OOM
现象:日志报java.lang.OutOfMemoryError,Task 失败重试。原因:单个分区数据过大,或collect()把大结果拉回 Driver。解决:调大spark.executor.memory,同时调大spark.sql.shuffle.partitions让分区更细;杜绝在 Driver 端collect大表,改用write落盘。
5.3 时间戳时区错乱
现象:按天聚合的结果和业务方对不上,差几个小时。原因:from_unixtime默认用集群时区,集群配的是 UTC 而业务要东八区。解决:在from_unixtime里显式传时区,或启动时设spark.sql.session.timeZone=Asia/Shanghai,统一口径。
5.4 shuffle partitions 默认 200 拖慢小任务
现象:本地跑几万条数据也要几十秒,日志里 200 个 Task 大部分空跑。原因:spark.sql.shuffle.partitions默认 200,小数据量下调度开销大于计算。解决:本地调试设成 8 或 16,集群按数据量设成核数的 2 到 3 倍。
5.5 会话切分阈值拍脑袋定
现象:人均会话数异常高或异常低,业务方质疑。原因:30 分钟阈值不适用当前场景,或时间戳单位是毫秒被当成秒。解决:先确认时间戳单位(秒还是毫秒),再和业务方确认会话定义,阈值写进配置而不是硬编码。
6. 验证与进阶:怎么确认这套分析系统真的算对了
跑通不等于算对。验证分三层:数据层、逻辑层、业务层。
数据层验证:清洗前后条数差、去重前后差、空值率,这些用count和filter就能查。我习惯在每步清洗后打一条日志,记录输入输出条数,出问题时能快速定位是哪一步吃掉了数据。
逻辑层验证:拿一小批已知答案的数据手工算一遍,和 Spark 结果对比。比如手工数 100 条日志里的购买用户数,和 SQL 结果核对。漏斗转化率、留存率这类指标,用小数据集验证公式没写反。
业务层验证:和业务方已有的报表对,差异超过 5% 就要查口径。常见差异来源是时区、去重规则、会话阈值。
进阶方向有两个。一是把批处理换成 Structured Streaming,做实时漏斗和实时大屏,核心代码逻辑不变,把read换成readStream、write换成writeStream即可,但要注意 watermark 和状态管理。二是把 RFM 分层结果写回 Hive 或 ClickHouse,供运营系统调用,形成"分析 → 打标 → 触达"的闭环。
# 流式版本的核心改动:从批到流的三个替换 stream_df = spark.readStream.schema(schema).json("hdfs:///data/user_behavior/") \ .withWatermark("event_time", "10 minutes") # 容忍 10 分钟乱序 query = stream_df.writeStream \ .outputMode("append") \ .format("parquet") \ .option("checkpointLocation", "hdfs:///checkpoint/funnel") \ .start("hdfs:///result/funnel")逻辑说明:withWatermark定义乱序容忍度,超过 watermark 的迟到数据会被丢弃;checkpointLocation必须指定,否则重启后状态丢失。流式漏斗的难点在状态管理,会话跨批次时要靠flatMapGroupsWithState维护,比批处理复杂一个量级,建议批处理稳定后再上。
参数说明:watermark 设太小会丢迟到数据,设太大会增加状态内存。10 分钟是电商场景的常见起点,按实际乱序情况调。
最后说个习惯:这套系统我踩过最大的坑不是代码,是口径。同一份数据,运营要的"活跃用户"和产品要的"活跃用户"可能差 20%。所以每次动手前,先把指标定义写成文档,和需求方确认签字,再写代码。源码可以复用,口径不能想当然。希望帮到你。
本文还有配套的精品资源,点击获取