简介:这份Spark电商用户行为分析系统源码与完整文档,面向计算机相关专业的毕业生、课程设计参与者及希望积累大数据实战经验的学习者,可直接用于毕业设计、课程作业或项目实践。系统基于Spark分布式计算架构,支持实时分析与离线计算两种模式,涵盖用户点击流分析、购买行为模式识别、用户画像构建等核心功能,并整合Spark MLlib协同过滤算法实现个性化推荐,借助Spark Streaming处理实时数据流,可视化部分采用ECharts展示分析结果。资源包共286个文件,以185个xml配置、40个java源码、47个zbak备份及少量properties、png、md等文档为主,压缩包约1.28MB,目录结构清晰,便于按模块检索学习。目前已有56人学习下载。资料内含全部必要技术文档与配置说明,下载后即可部署运行,是学习大数据技术、理解电商用户行为分析完整链路的优质实践案例。
1. 从一份能跑通的 Spark 电商用户行为分析系统说起
电商后台每天滚出几千万条点击、加购、下单日志,老板要的是「昨天有多少人从浏览走到支付」「哪个品类跳失最狠」,而数据团队手里只有一堆 JSON 文本和一台还没配好的集群。这份 Spark 电商用户行为分析系统源码与完整文档,就是冲着这个场景来的:它把数据清洗、会话切分、漏斗转化、品类排行这几条主线用 Spark 串起来,配了完整文档说明每个模块的输入输出。适合正在做 spark 毕业设计的学生,也适合刚接手用户行为埋点、想找一份能对照着改的工程骨架的初中级数据开发。源码不是玩具 demo,它按真实日志字段设计,跑通之后你能直接换成自己公司的埋点格式。
2. 环境与数据准备:集群搭建和日志字段对齐
2.1 为什么选 Spark 而不是单机 Pandas
用户行为日志的量级很尴尬:几十 GB 用 Pandas 内存扛不住,上 Flink 又嫌重。Spark 的 DataFrame API 对这类「读日志、做聚合、写结果」的批处理任务刚好合适,而且 spark sql 的窗口函数能直接表达会话切分和漏斗步骤,不用自己写状态机。这份源码用的是 Spark SQL + DataFrame 混合写法,核心逻辑集中在几个 SQL 里,改起来比纯 RDD 直观。选 Spark 还有一个现实原因:spark 集群搭建的资料多,头歌平台上就有 spark 的安装与使用实验,学生照着配环境不至于卡在第一步。
2.2 集群搭建的最小可用配置
单机伪分布式足够跑通这份源码,生产再换 standalone 或 YARN。下面是我一般会用的最小步骤,基于 Linux,JDK 8 或 11 都行。
# 下载并解压 Spark(版本按你文档里写的来,这里以 3.x 为例) tar -zxvf spark-3.x.x-bin-hadoop3.tgz -C /opt/ cd /opt/spark-3.x.x-bin-hadoop3 # 配置环境变量,追加到 ~/.bashrc export SPARK_HOME=/opt/spark-3.x.x-bin-hadoop3 export PATH=$SPARK_HOME/bin:$PATH export JAVA_HOME=/usr/lib/jvm/java-11-openjdk # 伪分布式需要指定 master,本地模式可直接 local[*] # 验证安装 spark-submit --version逻辑说明:SPARK_HOME让 spark-submit 能找到依赖,local[*]表示用本机所有核跑,调试阶段够用。参数上,spark.executor.memory在伪分布式里就是 JVM 堆大小,给 2g 起步;spark.sql.shuffle.partitions默认 200,小数据量下会拖慢任务,源码文档里如果没提,你手动改成 8 或 16,能明显减少小文件。
2.3 日志字段对齐:别急着跑,先看列名
这份源码假设的原始日志字段大致是:user_id、item_id、category_id、behavior_type(pv/cart/fav/buy)、timestamp。你拿到的埋点数据列名大概率不一样,常见做法是在读取后立刻做一次 rename,而不是改 SQL。
from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_timestamp spark = SparkSession.builder \ .appName("ecommerce_behavior") \ .config("spark.sql.shuffle.partitions", "16") \ .getOrCreate() # 读取原始 JSON 日志,字段名按你实际埋点调整 raw = spark.read.json("hdfs:///data/user_behavior/*.json") # 统一列名,避免后续 SQL 到处改 df = raw.select( col("uid").alias("user_id"), col("item").alias("item_id"), col("cate").alias("category_id"), col("bhv").alias("behavior_type"), to_timestamp(col("ts"), "yyyy-MM-dd HH:mm:ss").alias("event_time") ) df.createOrReplaceTempView("behavior")逻辑说明:to_timestamp把字符串时间转成 Spark 的 TimestampType,后面做会话切分和时间窗口全靠它。参数上,格式串必须和日志里的时间格式完全一致,差一个字符就全变 null,这是最常见的翻车点。createOrReplaceTempView注册成临时视图后,源码里那些 SQL 就能直接跑。
提示:先
df.printSchema()和df.show(5)确认字段类型,再往下走。时间列如果是 long 型毫秒戳,用(col("ts")/1000).cast("timestamp")。
3. 核心分析模块拆解:会话切分、漏斗与品类排行
3.1 会话切分:用窗口函数替代状态机
用户行为分析里「一次会话」通常定义为:同一用户,相邻行为间隔超过 30 分钟就算新会话。源码里用 Spark SQL 的lag窗口函数实现,比写 UDF 状态机清爽得多。
-- 计算每条行为与上一条行为的时间差,超过 1800 秒则标记为新会话 WITH ordered AS ( SELECT user_id, item_id, behavior_type, event_time, LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time) AS prev_time FROM behavior ), flagged AS ( SELECT *, CASE WHEN prev_time IS NULL THEN 1 WHEN UNIX_TIMESTAMP(event_time) - UNIX_TIMESTAMP(prev_time) > 1800 THEN 1 ELSE 0 END AS is_new_session FROM ordered ) SELECT user_id, item_id, behavior_type, event_time, SUM(is_new_session) OVER (PARTITION BY user_id ORDER BY event_time) AS session_id FROM flagged逻辑说明:第一层ordered用LAG拿到同用户上一条行为时间;第二层判断间隔是否超阈值,是则打 1;第三层用累加和给每条行为分配会话编号。参数上,1800 秒是电商场景常用值,短视频或资讯类可能缩到 300 秒,改这个数就行。注意PARTITION BY user_id不能少,否则会把不同用户的行为串在一起。
3.2 转化漏斗:pv → cart → buy 的三步统计
漏斗是这份源码最实用的部分,它按会话粒度统计每个用户是否走完三步,再算整体转化率。
-- 按会话聚合,标记是否发生过各行为 WITH session_agg AS ( SELECT user_id, session_id, MAX(CASE WHEN behavior_type = 'pv' THEN 1 ELSE 0 END) AS has_pv, MAX(CASE WHEN behavior_type = 'cart' THEN 1 ELSE 0 END) AS has_cart, MAX(CASE WHEN behavior_type = 'buy' THEN 1 ELSE 0 END) AS has_buy FROM session_behavior GROUP BY user_id, session_id ) SELECT SUM(has_pv) AS pv_sessions, SUM(has_cart) AS cart_sessions, SUM(has_buy) AS buy_sessions, ROUND(SUM(has_cart) / SUM(has_pv), 4) AS pv_to_cart_rate, ROUND(SUM(has_buy) / SUM(has_cart), 4) AS cart_to_buy_rate FROM session_agg逻辑说明:MAX(CASE WHEN ...)是典型的行转列技巧,把同一会话内的多种行为压成一行标记。参数上,转化率用ROUND保留四位小数,避免除零可以加NULLIF。这份源码的文档里如果写了「漏斗按用户去重」,那GROUP BY就换成 user_id,两种口径结果差很多,看业务定义。
3.3 品类排行与 spark sql 日期处理
品类维度的排行通常按 GMV 或订单数排,源码里用的是订单数。这里顺带说一个高频需求:按月份聚合。spark sql 日期转存月份用date_format或trunc。
-- 按品类统计购买次数,并按月份分组 SELECT category_id, DATE_FORMAT(event_time, 'yyyy-MM') AS month, COUNT(*) AS buy_count FROM behavior WHERE behavior_type = 'buy' GROUP BY category_id, DATE_FORMAT(event_time, 'yyyy-MM') ORDER BY buy_count DESC逻辑说明:DATE_FORMAT把时间戳转成yyyy-MM字符串,适合做月度报表。如果要做日期加减,用DATE_ADD(event_time, 7)或ADD_MONTHS。参数上,ORDER BY在数据量大时建议配合LIMIT,否则全量排序会拖慢。这份源码的品类排行模块还加了RANK()窗口函数,能直接出 Top N。
注意:
spark.sql.shuffle.partitions在聚合和排序时影响巨大,小数据集下 200 个分区会产生大量空任务,改成核数的 2~4 倍。
4. 避坑与排查:那些让我重跑过任务的坑
4.1 时间格式不匹配导致全表 null
现象:会话切分结果全是 0,prev_time全为 null。原因:to_timestamp的格式串和日志实际格式不一致,比如日志是2024/01/01 10:00:00,代码写的是yyyy-MM-dd。解决:先df.select("ts").show(5, false)看原始值,再对齐格式串,或者干脆用from_unixtime处理毫秒戳。
4.2 shuffle 分区过多拖慢小数据任务
现象:本地跑 10 万条数据要几分钟,日志里全是ShuffleMapTask。原因:默认 200 个 shuffle 分区,每个分区只有几百条记录,调度开销远大于计算。解决:在 SparkSession 里设spark.sql.shuffle.partitions=16,或者按spark.default.parallelism调整。
4.3 left outer join 广播右侧的误用
现象:任务卡在 join 阶段,或者报 OOM。原因:spark 对 left outer join 只能广播右侧,如果右侧表很大,广播会撑爆 driver。解决:确认小表在右,或者显式用/*+ BROADCAST(t2) */提示,大表 join 就老老实实走 shuffle。
4.4 会话切分漏掉 PARTITION BY
现象:不同用户的行为被算进同一个会话,session_id 串号。原因:窗口函数里只写了ORDER BY event_time,没写PARTITION BY user_id。解决:所有按用户维度的窗口操作,PARTITION BY user_id必须加,这是血泪经验。
4.5 输出小文件过多
现象:结果目录下几百个几百字节的文件。原因:分区数多 + 每个分区单独写。解决:写出前repartition(1)或coalesce(1),但注意 coalesce 不触发 shuffle,数据量大时可能单分区内存不够,用 repartition 更稳。
5. 进阶技巧:把这份源码改成你自己的分析管道
跑通默认流程只是开始,真正省事的是把它变成可配置的管道。我一般会做三件事:把会话超时、漏斗步骤、输出路径抽成配置文件;把 SQL 从代码里挪到独立.sql文件用spark.sql(open(...).read())加载;加一层数据质量校验,比如空 user_id 比例超过 1% 就告警。验证方法很简单:拿一份已知结果的小数据集跑一遍,对比手算的转化率,误差在千分位以内就算对。另一个技巧是 spark 内存线程监测,在spark-submit时加--conf spark.ui.port=4040,通过 Web UI 看每个 stage 的 shuffle 读写和 GC 时间,比看日志快得多。从那以后我每次改完 SQL 都强制走一遍小数据集回归,确认漏斗数字没跳变才上全量。希望帮到你。
本文还有配套的精品资源,点击获取