
简介基于Spark的电商用户行为分析系统完整项目源码面向计算机类毕业设计、课程设计及Spark入门学习者。资源围绕用户会话分析、商品Top3统计等典型场景展开代码采用Scala编写结构清晰代码模块划分明确覆盖数据加载、清洗、统计分析与结果输出等环节可帮助读者快速掌握Spark核心算子应用与电商数据分析思路。压缩包共273个文件其中208个xml多用于项目配置与数据定义58个scala文件为源码主体另有properties配置文件、gitignore版本管理文件、iml工程文件及md说明文档整体仅183KB轻量易用。目前已有125人浏览学习具备一定参考价值。读者可从该资源中获取完整项目骨架、配置方式与业务功能实现适合作为毕业设计参考或课程设计练手项目无论是理解项目工程结构还是复用算法逻辑都有实际帮助。1. 这个 Spark 项目解决的远不止“用户行为分析”做电商数据分析的同学应该都有体会用户在平台上的每一次点击、搜索、下单单独看都只是日志里的一行记录但把它们按 session 切分、按用户聚合、按区域统计之后就能回答“用户从进站到下单经历了什么”“哪些商品在哪些区域卖得好”“整个转化漏斗在哪一步流失最严重”这类业务问题。这个基于 Spark 的电商用户行为分析系统正是围绕这三类问题展开的——它不是简单的“读数据-出报表”demo而是一套包含用户访问 session 聚合统计、页面单跳转化率计算、区域热门商品 Top3 挖掘的完整离线分析工程。项目源码基于 Scala 编写核心文件是Demand1Function.scala、UserSessionAnalysisFunction2.scala和AreaTop3ProductFunc.scala配合commerce.properties配置文件管理输入输出路径和参数。从文件命名就能看出这是按“需求 1、需求 2、需求 3”的课程设计惯例组织的但代码结构和实现方式基本遵循了生产环境 Spark 离线任务的标准写法。适合三类人准备做 Spark 方向毕业设计的学生想补全 Spark SQL DataFrame 实战经验的初级开发以及需要一套可扩展的分析模板、想快速改造接入自己业务的工程师。有一点需要先说清这个项目用的技术栈是 Spark Core Spark SQL核心 API 是 DataFrame 和 Dataset而不是纯 RDD 编程。这是个务实的选择——在 Spark 2.x 时代DataFrame 的 Catalyst 优化器会在你写代码时就帮你做谓词下推和列剪枝同样的分析逻辑用 RDD 写不仅代码量大执行效率也未必更好。理解了这一点你才能真正看懂这个项目为什么能扛住“千万级用户行为数据”这类课程设计里最常见的性能要求。下面从数据落地开始逐步把三个需求的实现拆开讲。2. 数据落地与 user_visit_action 的预处理逻辑2.1 从原始日志到可分析的宽表电商用户行为分析的第一步不是写分析代码而是把日志数据整理成一张结构清晰、字段完备的明细表。这个项目里的核心数据表是user_visit_action它承载了用户每一次行为动作的完整记录。典型字段包括session_id会话 ID、user_id用户 ID、action_time行为时间、search_keyword搜索关键词、click_category_id点击的商品分类 ID、click_product_id点击的商品 ID、order_category_ids下单的商品分类 ID、order_product_ids下单的商品 ID、pay_category_ids支付的分类 ID、pay_product_ids支付的商品 ID以及city_id城市 ID。在离线分析场景里这张表通常由 Hive 管理数据以 Parquet 或 ORC 格式存储在 HDFS 上。Parquet 是列式存储格式对于“只读取其中几个字段做聚合”的分析任务它能跳过无关列I/O 开销大幅低于行式存储。如果你用 CSV 或 JSON 作为数据源建议在项目开始时先做一次格式转换——这个动作可以用一条简单的 Spark SQL 完成也可以在建表时直接指定STORED AS PARQUET。CREATE TABLE user_visit_action ( session_id STRING, user_id BIGINT, action_time TIMESTAMP, search_keyword STRING, click_category_id BIGINT, click_product_id BIGINT, order_category_ids STRING, order_product_ids STRING, pay_category_ids STRING, pay_product_ids STRING, city_id BIGINT ) PARTITIONED BY (dt STRING) STORED AS PARQUET;分区字段dt按日期组织数据这一步很关键。离线分析任务通常是 T1 调度每天处理前一天的全量或增量数据按天分区能让你在读取数据时直接做分区裁剪——只扫描需要的日期而不是把整张表从头读一遍。实际项目中上游数仓任务会在每天凌晨把前一天的日志写入对应分区分析任务在早上执行时直接WHERE dt 2025-01-15即可。2.2 Spark Session 初始化与配置参数项目入口使用 SparkSession 统一入口替代了旧的 SparkContext SQLContext 组合。commerce.properties文件里管理了运行时的关键参数包括输入表名、输出路径、session 超时阈值等这样代码里不写死业务参数换一套数据或调整分析窗口时只改配置文件即可。val spark SparkSession.builder() .appName(EcommerceUserBehaviorAnalysis) .config(spark.sql.shuffle.partitions, 200) .config(spark.sql.adaptive.enabled, true) .config(spark.sql.adaptive.coalescePartitions.enabled, true) .enableHiveSupport() .getOrCreate()spark.sql.shuffle.partitions控制的是 shuffle 后的默认分区数默认值 200但在数据量只有几 GB 的测试环境里200 个分区会产生大量小文件反而拖慢后续读取。我一般在课程设计或数据量级在亿级以下的任务里把它调到 50100 之间真实集群上则要按总数据量 / 每个分区目标大小约 128MB256MB来估算。spark.sql.adaptive.enabled是 Spark 3.0 引入的 Adaptive Query Execution它能在运行时根据中间结果的实际大小动态调整 shuffle 分区数减少小文件问题也缓解数据倾斜。如果你的 Spark 版本低于 3.0这个配置不生效需要手动调分区数。enableHiveSupport()让 Spark 可以读写 Hive 表。如果你不想依赖 Hive 元数据也可以直接从 HDFS 路径读 Parquet 文件但那样就需要自己在代码里维护 schema数据字段一多就很容易出错。2.3 数据质量检查空值和脏数据的处理策略数据落地后不能急着开跑需要先做一轮质量检查。用户行为日志经常有这些问题session_id 为空埋点丢失、action_time 不在合理时间范围时钟偏移、click_product_id 和 order_product_ids 同时为空页面浏览行为没有产生点击和下单。SELECT COUNT(*) AS total_cnt, SUM(CASE WHEN session_id IS NULL OR session_id THEN 1 ELSE 0 END) AS null_session_cnt, SUM(CASE WHEN action_time 2025-01-15 00:00:00 OR action_time 2025-01-15 23:59:59 THEN 1 ELSE 0 END) AS abnormal_time_cnt FROM user_visit_action WHERE dt 2025-01-15;这些检查结果决定了后续分析的过滤策略。对于 session_id 为空的数据因为后续所有的 session 聚合都依赖这个字段直接过滤掉是最稳妥的做法对于时间异常的数据如果数量占比很低低于 0.1%也可以直接丢弃但如果占比偏高就要复查埋点逻辑是否出了问题。这个判断在写分析代码之前必须完成否则算出来的 session 时长分布、转化率指标都会被脏数据污染。3. 用户访问 Session 聚合分析从基础统计到漏斗拆解3.1 按 Session 聚合的核心思路Session 是指用户在一次访问过程中产生的连续行为序列通常以“连续 30 分钟无新操作则视为会话结束”来切分。UserSessionAnalysisFunction2.scala处理的就是这个环节——它需要把user_visit_action中同一session_id下的所有行为汇总成一条 Session 记录统计出访问时长、步长、起始时间、结束时间等聚合指标。这里的实现方式很直接按session_id分组对action_time求 min 和 max 得到访问时间范围用count(*)得到步长用户在这个 Session 里产生了多少行为再关联user_info表拿到用户的年龄、性别、职业等维度信息为后续的分维度分析做准备。val sessionAggrDF userVisitActionDF .groupBy(session_id) .agg( min(action_time).as(start_time), max(action_time).as(end_time), count(*).as(step_length), collect_list(click_product_id).as(clicked_products) ) .join(userInfoDF, user_id)collect_list的作用是把 Session 内点击过的所有商品 ID 收集成一个数组这为后续“判断用户是否点击过某个商品”提供了便捷的查询方式。不过要注意的是如果 Session 内的行为数据量非常大比如一个用户的 Session 跨越了一整天collect_list可能会产生较大的序列化开销。常见的做法是只保留必要字段或者在聚合前先对明细数据做一轮筛选。3.2 访问时长与步长的分桶统计拿到 Session 聚合结果后业务上最关心的两个指标是“访问时长”和“访问步长”因为它们直接反映用户粘性和访问深度。项目的做法是对这两个指标做分桶统计——把访问时长按 03 秒、310 秒、1030 秒、3060 秒、13 分钟、3 分钟以上分桶把访问步长按 13 次、46 次、79 次、1030 次、30 次以上分桶然后统计每个桶内的 Session 数量以及占比。SELECT CASE WHEN (unix_timestamp(end_time) - unix_timestamp(start_time)) BETWEEN 0 AND 3 THEN 0-3s WHEN (unix_timestamp(end_time) - unix_timestamp(start_time)) BETWEEN 3 AND 10 THEN 3-10s WHEN (unix_timestamp(end_time) - unix_timestamp(start_time)) BETWEEN 10 AND 30 THEN 10-30s WHEN (unix_timestamp(end_time) - unix_timestamp(start_time)) BETWEEN 30 AND 60 THEN 30-60s WHEN (unix_timestamp(end_time) - unix_timestamp(start_time)) BETWEEN 60 AND 180 THEN 1-3m ELSE 3m END AS duration_bucket, COUNT(*) AS session_cnt FROM sessionAggrDF GROUP BY duration_bucket ORDER BY duration_bucket;这里用unix_timestamp把时间字段转成秒数再求差避免了对 TIMESTAMP 直接做减法时可能出现的类型不匹配问题。分桶边界值用BETWEEN的闭区间要特别小心——03 秒和 310 秒在 3 秒这个边界上会重复计数所以在设计桶时需要明确是左闭右开还是全闭区间并在多个桶之间避免重叠。如果你直接用原项目代码不检查这个边界最终的分布图在边界点上会出现一个异常的凸起。3.3 按用户维度聚合多 Session 合并与维度关联如果一个用户一天之内产生了多个 Session分析用户维度的行为时就需要把这些 Session 合并起来。常见做法是以user_id为粒度再次聚合访问总时长、总步长、Session 数量以及用户感兴趣的商品类别集合。val userAggrDF sessionAggrDF .groupBy(user_id) .agg( sum(step_length).as(total_step), count(session_id).as(session_cnt), collect_set(click_category_id).as(interest_categories) )collect_set在这里比collect_list更合适——它自动去重得到的商品类别集合是用户真正的兴趣范围不会因为用户反复点击同一类目而放大权重。聚合完成后关联user_info表中的年龄、性别、职业字段就能得到“2530 岁女性用户平均访问时长最长”“程序员群体在凌晨时段访问占比最高”之类的业务洞察。这段逻辑的关键在于理解聚合层级明细行为 → Session 聚合 → 用户聚合每一层都是一次groupBy。如果你在用户聚合层发现数据量异常膨胀多半是前面的 Session 切分出了问题比如 session_id 生成逻辑不规范导致多个用户共用同一个 session_id或者同一个用户一天内产生了数百个极短 Session。遇到这种情况先回查明细数据的 session_id 分布不要在下游调参。4. 页面单跳转化率剖析用户路径的每一步流失4.1 什么是页面单跳转化率页面单跳转化率衡量的是用户从页面 A 跳转到页面 B 的概率计算公式是A-B 的跳转次数 / A 页面的访问次数。在电商场景里这个指标常用于分析核心漏斗——比如“首页 → 搜索页 → 商品详情页 → 下单页 → 支付页”的每一跳转化率从而定位流失最严重的环节。这个需求在代码上比 Session 聚合复杂得多难的不是统计本身而是如何从原始行为数据中准确还原用户的页面访问序列。Demand1Function.scala处理的就是这个环节——它的做法是按 session 和时间排序把用户的行为序列切成形如page1 - page2 - page3的相邻页面跳转对然后按跳转对聚合计数。4.2 相邻页面跳转对的生成算法生成跳转对的核心代码涉及窗口函数和集合操作import org.apache.spark.sql.expressions.Window val pageSeqDF userVisitActionDF .withColumn( rn, row_number().over(Window.partitionBy(session_id).orderBy(action_time)) ) .filter(rn 2) .select( $session_id, $action_time, $click_category_id, lag(click_category_id, 1).over( Window.partitionBy(session_id).orderBy(action_time) ).as(prev_page) )这段代码用lag窗口函数取同一 Session 内上一行记录的页面 ID与当前行拼接后得到跳转对。为什么不用lead因为lead是往下取生成的跳转对会变成 “当前页 - 后一页”在计算“从 A 到 B 的转化率”时你需要的是以 A 为起点的所有跳转去向而不是以 B 为终点的所有来源所以lag更直观。filter(rn 2)的目的是去掉每个 Session 的第一条记录——它没有上一页无法构成跳转对。这里有个性能陷阱必须提醒lag配合orderBy action_time时Spark 会对每个分区内的数据进行全排序如果某个 Session 的行为记录特别多达到数万条这个排序的开销会比较大。实际项目中我会在窗口排序前先按action_time做一次全局重分区然后再在分区内部用sortWithinPartitions排好序这样窗口函数执行时的排序压力会小很多。另一种常见做法是去掉row_number这步直接对lag的结果做WHERE prev_page IS NOT NULL过滤效果等价但少一次窗口计算。4.3 转化率计算与多级漏斗的 SQL 实现生成跳转对之后就可以按跳转对做聚合统计了。下面是一段基于 Spark SQL 的计算脚本WITH page_jump_count AS ( SELECT prev_page AS source_page, click_category_id AS target_page, COUNT(*) AS jump_cnt FROM page_jump_pair WHERE prev_page IS NOT NULL GROUP BY prev_page, click_category_id ), page_visit_count AS ( SELECT click_category_id AS page_id, COUNT(DISTINCT session_id) AS visit_cnt FROM page_jump_pair GROUP BY click_category_id ) SELECT j.source_page, j.target_page, j.jump_cnt, v.visit_cnt, ROUND(j.jump_cnt / v.visit_cnt, 4) AS conversion_rate FROM page_jump_count j JOIN page_visit_count v ON j.source_page v.page_id ORDER BY j.source_page, conversion_rate DESC;注意这里的visit_cnt统计维度——用COUNT(DISTINCT session_id)而不是COUNT(*)因为一个 Session 内可能多次访问同一个页面如果不按 Session 去重转化率会被高频用户的重复行为放大失真。电商漏斗的分析场景里我们关心的是“有多少用户走到了这一步”而不是“这些页面被访问了多少次”。如果你的业务需要分析的是完整漏斗比如首页 → 商品页 → 下单页 → 支付页可以在得到跳转对表之后用filter把跳转对限制在漏斗路径内的组合然后逐级计算转化率。项目代码里没有显式实现漏斗路径的配置化但你可以通过commerce.properties传入一个funnel.pages参数来实现这样不同业务线可以复用同一套计算逻辑。5. 区域热门商品 Top3维度下钻与窗口函数的配合5.1 区域与商品的关联维度设计第三个需求是统计每个区域通常按省份或城市粒度点击量最高的前 3 个商品。这个需求比前两个多了一层维度的复杂性——原始数据里只有city_id没有区域名称需要关联城市维度表才能拿到省份信息。同时商品点击数要从行为表中提取而商品本身的名称、价格等信息又存在商品表中所以这是一个典型的多表关联 分组 TopN 问题。val cityInfoDF spark.read .format(jdbc) .option(url, jdbcUrl) .option(dbtable, city_info) .load() val productInfoDF spark.read .format(jdbc) .option(url, jdbcUrl) .option(dbtable, product_info) .load()在实际项目中城市维度和商品维度通常存放在 MySQL 或 PostgreSQL 里因为它们是变化缓慢的维度数据用关系型数据库管理更合适。行为数据量太大不适合放在关系型数据库里所以留在 Hive/数仓中。但用 JDBC 直连 MySQL 的方式读取维度表在数据量大时会有性能隐患——city_info和product_info通常只有几百到几万行全量加载到 Spark 内存中做广播 join 是可行的但如果维度表增长到千万级必须改用 Hive 表或分布式缓存。5.2 多表 Join 与区域商品点击量的聚合先把行为数据、城市数据、商品数据关联起来得到一张包含“省份、城市、商品、点击量”的明细宽表。这里的重点是 join 的先后顺序——先用小表城市信息做广播 join再用大表商品信息做 shuffle join这样能显著减少 shuffle 的数据量。SELECT ci.province_name, uva.city_id, uva.click_product_id, pi.product_name, COUNT(*) AS click_cnt FROM user_visit_action uva JOIN city_info ci ON uva.city_id ci.city_id JOIN product_info pi ON uva.click_product_id pi.product_id WHERE uva.click_product_id IS NOT NULL AND uva.click_product_id ! -1 GROUP BY ci.province_name, uva.city_id, uva.click_product_id, pi.product_name过滤条件click_product_id ! -1很重要因为很多埋点系统用 -1 表示“没有点击商品”的行为。如果不加这个过滤-1 会作为一个虚假的商品 ID 参与聚合并且大概率因为商品表中找不到对应记录而在 join 时被丢弃但这样会导致计数不一致——行为表里统计到 100 万条点击join 后只剩 80 万条数据对不上。5.3 窗口函数 row_number 实现分组 TopN分组 TopN 是窗口函数的经典场景——按省份分组组内按点击量降序排列然后取前 3 名SELECT province_name, product_id, product_name, click_cnt FROM ( SELECT province_name, click_product_id AS product_id, product_name, click_cnt, row_number() OVER ( PARTITION BY province_name ORDER BY click_cnt DESC, product_id ASC ) AS rn FROM regional_product_click ) t WHERE rn 3为什么用row_number而不是rank或dense_rank因为这里的需求是“取 Top3 商品”如果两个商品的点击量并列第 3rank会返回 3 个以上的商品比如 4 个不符合“Top3”的严格定义。row_number会给每条记录一个唯一且连续的编号即使点击量相同也会按product_id的排序分出先后保证每个区域最多返回 3 条记录。这段代码是整个项目的点睛之笔——它展示了一个非常重要的统计思维离线分析里“TopN”通常是有歧义的需求必须明确你期望的是“最多取 N 个”用row_number还是“所有达到第 N 名成绩的都要”用rank。业务方如果没想清楚你在实现前必须先问清楚否则返工成本很高。6. 调优与验证让这个分析任务在真实数据上跑得更稳6.1 数据倾斜的定位与三种常用缓解手段三个需求的计算逻辑都涉及groupBy或join数据倾斜是跑真实数据时大概率会遇到的问题。倾斜的症状是某个 reducer 长时间运行、其他 reducer 早已完成整个 Spark 任务卡在最后一个 stage。定位方法很简单——在 Spark UI 的 Stages 页面查看每个 task 的 Shuffle Read 大小如果有 task 处理的数据量是其他 task 的 10 倍以上基本可以确认倾斜。缓解手段按优先级排有三种第一检查是否能通过过滤异常 key 来消除倾斜比如上面提到的 -1 商品 ID它在 groupBy 时会造成单 key 数据量巨大直接过滤掉即可第二对热点 key 加随机前缀后再聚合分两步做——先按key 随机数聚合一次去掉随机数后再聚合一次这个方法对 groupBy 倾斜非常有效对标count、sum这类聚合函数不需要修正但对avg需要先把总和和计数分别算出来再除第三如果是 join 倾斜考虑把热点 key 的数据单独拆出来广播非热点 key 走常规 shuffle join两者 union 得到最终结果。6.2 合理配置 Spark 执行参数课程设计用的数据量通常不大但你可能想把项目跑出“看起来很专业”的性能那么这几个参数值得关注。spark.executor.memory建议按每个 executor 上的并发 task 数来规划——如果一个节点有 16GB 内存、跑 4 个 executor每个 executor 分配 3GB 就够用留出余量给系统和其他进程。spark.executor.cores通常设 24设太高会导致 executor 间磁盘和网络的 I/O 竞争反而降低整体吞吐。spark.default.parallelism可以根据总核心数 × 23来设定。比如集群有 4 个节点、每个节点 8 核总核心数 32那么可以用spark.default.parallelism64这样每个核心平均分配到 2 个任务既不会因为任务太少导致资源闲置也不会因为任务太多增加调度开销。spark.sql.shuffle.partitions在 3.0 以上版本可以配合 AQE 使用开启之后 Spark 会在 shuffle 结束时根据分区大小自动合并小分区。6.3 结果验证方法从“跑通”到“跑对”代码跑通只是第一步结果正确与否需要验证。推荐的做法是把结果拆成两个层面核验。第一层是总量校验。对每个需求的输出和明细数据算一个对账指标。比如 session 聚合结果中的 session 总数应该和明细表中SELECT COUNT(DISTINCT session_id)的结果一致。如果少了说明groupBy前有数据被意外过滤如果多了说明 session_id 有重复或空值没有正确处理。第二层是抽样核对。从最终结果中随机抽取 23 个区域或 session回到明细数据里手工计算一遍和输出结果对比。比如验证某个省份的商品 Top3可以手工执行一条 SQLSELECT click_product_id, COUNT(*) AS cnt FROM user_visit_action uva JOIN city_info ci ON uva.city_id ci.city_id WHERE ci.province_name 广东省 AND uva.dt 2025-01-15 AND uva.click_product_id IS NOT NULL GROUP BY uva.click_product_id ORDER BY cnt DESC LIMIT 3;把这条 SQL 的结果和程序输出的结果对比。不要只看前 3 名是否一致还要看点击量数值是否对得上——如果数量差了检查是不是 join 时丢了数据或者窗口函数的排序字段选得不对。还有一种常见错误ORDER BY click_cnt DESC的时候如果 click_cnt 相同MySQL/Spark 的默认排序不保证稳定导致两次运行结果不同。要在排序字段后面追加一个唯一字段比如product_id作为次级排序条件才能保证结果可复现。本文还有配套的精品资源点击获取