Spark电商用户行为分析:RDD统计PV/UV并写入MySQL 📅 发布时间:2026/9/16 19:27:48 👁 浏览次数: 简介基于Spark电商用户行为分析的代码程序包面向大数据初学者与Spark开发者围绕用户浏览、点击、下单、支付等行为数据进行完整分析可直接作为课程设计或业务分析的基础源码。压缩包内共十七个文件以八份采用Scala语言编写的源文件为主体配合四份txt测试数据、SQL建表脚本、Maven依赖配置、properties配置、csv样例与Markdown格式说明等辅助文件包体大小约三点五三兆字节结构紧凑。目前已有二百四十四人学习下载适合希望快速获得可运行示例的读者。代码结合测试数据可立即验证页面浏览量、点击量等统计逻辑SQL脚本便于准备MySQL环境数据格式文档则降低了上手门槛能够帮助使用者省去重复造轮子的时间聚焦核心业务分析有效提升开发效率。1. 为什么这份Spark电商用户行为分析代码值得直接拿来做基线很多Spark入门教程停留在WordCount但真正要处理用户行为数据时字段解析、动作类型过滤、结果落库这些环节才是耗时大头。这份源码给的不是抽象示例而是一套带业务语义的完整链路user_visit_action.txt记录用户每次访问动作product_info和city_info提供维度数据最终通过Spark作业输出PV、点击、下单、支付四类指标。适合已经能跑通官方示例、但对业务字段映射和MySQL结果写入不熟的人。直接拿现成代码改比从空文件开始省一晚上调试时间。尤其它把数据格式说明和建表语句一起打包能让你少踩不少脏数据解析的坑。2. 项目结构与数据格式先看明白输入才能改得动Spark作业2.1 源码目录里每个文件是干什么的从压缩包内文件列表看这是一个标准的Maven工程。pom.xml锁定Spark依赖src/main放正式逻辑src/test放测试代码testData目录存放user_visit_action.txt、product_info.txt、city_info.txt和user_visit_action_test.txt。这种布局在真实项目里很常见测试数据单独放不直接混入生产环境但每条记录的分隔符和字段顺序必须提前弄清楚。spark-analysis ├── pom.xml ├── src/main ├── src/test ├── testData │ ├── user_visit_action.txt │ ├── product_info.txt │ ├── city_info.txt │ ├── user_visit_action_test.txt │ └── user_visit_action.csv ├── sql │ └── mysql建表语句.sql └── 数据格式说明.md动手之前先把数据格式说明.md打开里面写了每个文件的字段含义。不要急着读代码因为Spark作业的map函数里字段下标一旦写错后面所有聚合结果都会错。看格式文档后再对照代码里的split分割逻辑通常一次就能定位字段位置。user_visit_action.csv要特别留意扩展名是csv不代表一定用逗号分隔跑之前先head几行确认。2.2 用户行为数据的字段顺序与解析要点user_visit_action.txt是核心输入每一行代表一次用户行为。常见做法是各字段用制表符切分具体顺序可以参照数据格式说明。为了演示我列一个典型的字段约定字段位置字段含义示例0日期2024-05-011用户IDuser_10012session IDsession_0013页面IDpage_0034动作类型0,1,2,35商品IDproduct_0016城市IDcity_001代码里如果直接用line.split(\t)拿到的是一个Array[String]下标取错会抛ArrayIndexOutOfBoundsException而不是安静地返回null。我一般做数据清洗时会先做一次字段长度校验把不合格的行写入日志而不是直接filter掉这样能保留排错现场。动作类型的枚举值也要核对0通常代表点击1代表下单2代表支付3可能是收藏或加入购物车。以源码中的实际定义为准。2.3 MySQL建表语句与Spark结果落库sql/mysql建表语句.sql提供了目标表结构。典型设计会包含日期、指标类型、数值、统计时间等字段主键往往设在stat_date和action_type上防止重复跑数时插入重复数据。Spark作业跑完后通过JDBC写入MySQL批量插入的batch size和连接参数是重点。CREATE TABLE user_behavior_stats ( stat_date VARCHAR(20), action_type VARCHAR(20), cnt BIGINT, create_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (stat_date, action_type) );val statsRDD: RDD[(String, String, Long)] ... statsRDD.foreachPartition { partition val conn DriverManager.getConnection(url, user, password) val stmt conn.prepareStatement( INSERT INTO user_behavior_stats(stat_date, action_type, cnt) VALUES (?, ?, ?) ) partition.foreach { case (date, actionType, cnt) stmt.setString(1, date) stmt.setString(2, actionType) stmt.setLong(3, cnt) stmt.addBatch() } stmt.executeBatch() stmt.close() conn.close() }逻辑说明这里没有用foreach逐条写是因为每个partition只会建立一次JDBC连接显著减少MySQL连接开销。addBatch攒够一批再executeBatch能把写入吞吐量提上去。参数说明JDBC URL 要带useServerPrepStmtstrue和rewriteBatchedStatementstrue才能让MySQL服务端真正支持批量执行否则驱动可能还是逐条插入性能提升有限。3. 从RDD到业务指标PV、点击、下单、支付四类统计的实现3.1 为什么源码选择RDD而不是DataFrame这份源码大量使用RDD算子。现在新项目都推DataFrame但RDD对字段下标的依赖更直观也更容易理解每一步map和reduce在干什么。对于数据量在千万级以内的用户行为分析RDD和DataFrame的性能差距不明显而且调试时可以在filter后加一个take(10)直接在driver端看数据形状。val rdd sc.textFile(testData/user_visit_action.txt) .map(_.split(\t)) .filter(_.length 7)逻辑说明textFile逐行读入split切分后得到字符串数组filter把字段不够的行挡在外面。参数说明这里length 7必须和实际字段数匹配。如果数据格式说明里写了15个字段写7只会把解析错误延后到取下标时。最稳妥的做法是_.length expectedFieldCount但遇到个别行本身缺失字段时需要先决定是丢弃还是补齐。3.2 用户PV统计一次reduceByKey完成PV就是用户每次页面访问的累加不去重。源码里通常取日期字段再对每个日期计数。val pvRDD rdd .map(line (line(0), 1)) .reduceByKey(_ _)逻辑说明reduceByKey会在分区内先合并相同key的计数再在shuffle端做第二次合并比groupByKey后map sum 少传很多数据。参数说明line(0)是日期下标如果你要按天统计就保留这个位置。如果还要按小时可以把下标改成小时字段所在的位置或者用line(0) _ line(1)作为复合key。如果业务指标要求按用户去重后的访问人数那就是UV不能用这个写法。PV和UV口径不要混运营报表里两者经常一起出现但Spark实现完全不同。3.3 点击、下单、支付动作类型字段过滤后聚合用户行为数据里动作类型通常用数字表示。源码里的做法是分别过滤出对应类型再按日期聚合。这里假设0是点击、1是下单、2是支付。val clickRDD rdd .filter(line line(4).equals(0)) .map(line (line(0), 1)) .reduceByKey(_ _) val orderRDD rdd .filter(line line(4).equals(1)) .map(line (line(0), 1)) .reduceByKey(_ _) val payRDD rdd .filter(line line(4).equals(2)) .map(line (line(0), 1)) .reduceByKey(_ _)逻辑说明三份RDD对应三个独立的Spark job每个job都会触发一次shuffle。如果测试数据量不大这是最简单直观的写法。参数说明line(4).equals(0)是字符串比较原始文件里的数字如果有空格或用了0\n需要用line(4).trim。动作类型字段下标请以数据格式说明为准不要照抄4。如果想优化可以只扫一遍数据用map(line ((line(0), line(4)), 1))做复合key聚合再用filter拆成多个结果RDD。这样能省两次全量扫描但代码可读性会降低。数据量到达亿级时再考虑这种优化。3.4 本地验证与Spark提交参数把三个RDD的结果保存到文件或MySQL。如果是本地验证用coalesce(1)合并成单文件再保存避免输出目录里出现几百个碎片文件。spark-submit \ --class com.example.UserBehaviorAnalysis \ --master local[2] \ --executor-memory 2g \ target/spark-analysis-1.0.jar \ testData/user_visit_action.txt testData/output逻辑说明local[2]用两个线程模拟并行执行方便在ide或命令行里调试不需要Yarn环境。参数说明--executor-memory 2g在处理全量user_visit_action.txt时如果不足改为4g或6g。--class必须写你实际主类的全限定名不要照抄否则会直接报ClassNotFoundException。打包时如果用mvn package默认打的jar可能不包含依赖。如果提交到Yarn集群需要把Spark依赖的jar一起带上或者使用maven-shade-plugin打fat jar。本地local[2]模式因为classpath里已经有Spark不打fat jar也能跑。3.5 结果验证与可靠性检查跑完后输出目录里会有part-00000这类文件。用cat testData/output/part-00000检查每行格式对照(日期,次数)是否符合预期。我一般会再做一个总行数校验用wc -l统计原始文件的行数再手动过滤点击类型算一下数量看看和Spark输出是否一致。wc -l testData/user_visit_action.txt awk -F \t $5 0 {count} END {print count} testData/user_visit_action.txt逻辑说明wc -l统计总访问行数awk按第5列过滤动作类型为0的行。参数说明-F \t指定制表符分隔如果真实数据是逗号要改成-F ,。这一步能快速发现split下标写错或过滤条件写反的问题不用反复提交Spark job。4. 测试数据与真实排错字段错位、Shuffle倾斜和MySQL连接溢出4.1 测试数据文件夹里有什么testData中user_visit_action_test.txt是专门切出来的小数据集适合用来跑通流程。user_visit_action.csv可能只是为了测试不同分隔符下的兼容性。product_info.txt和city_info.txt如果只统计PV用不上但做商品维度分析或城市维度分布时要通过cogroup或join关联进来。跑冒烟测试时先用_test.txt验证主流程再切全量数据。直接上来跑全量遇到字段错位时日志会被大量错误刷屏反而难定位问题。4.2 从打包到输出全流程在项目根目录执行下面命令。确保Maven配置了合适的仓库Spark和Scala的依赖能从中央仓库拉到。mvn clean package -DskipTests spark-submit --class com.example.UserBehaviorAnalysis --master local[2] \ target/spark-analysis-1.0.jar \ testData/user_visit_action_test.txt testData/output_small逻辑说明-DskipTests跳过测试避免测试类里如果引用了不存在的文件导致构建失败。参数说明输出目录testData/output_small要提前不存在否则Spark会报FileAlreadyExistsException。跑第二次前先删掉旧目录或者代码里用Files.deleteIfExists清理。如果日志里出现Exception in thread main java.lang.ArrayIndexOutOfBoundsException: 5代表取到了不存在的字段下标。这时直接用文本编辑器打开原始数据查看第几列确实是空值或者数据行之间分隔符不一致。user_visit_action.csv如果逗号分隔但代码写死了tab也会出现这种问题。4.3 常见故障对照表现象原因处理方式ArrayIndexOutOfBoundsException分隔符错误或字段数少于预期head查看真实行格式Executor Lost / OOM单个key数据量过大发生数据倾斜随机前缀加盐二次聚合MySQL too many connections每条记录都新建连接foreachPartition内复用连接输出文件数量过多没有合并分区coalesce(1) 或 adjust partition数结果偏大或偏小动作类型过滤条件写错用awk手动验证目标行数数据倾斜在电商用户行为里很常见比如某个热门商品被大量点击对应key的记录数远高于其他key。此时reduceByKey会把大量数据集中在同一个executor上。常见做法是给key加随机前缀先打散做第一轮聚合再去掉前缀做第二轮聚合。import scala.util.Random val saltedRDD rdd .map(line ((line(0) _ Random.nextInt(10), 1))) .reduceByKey(_ _) .map { case (key, cnt) val date key.split(_)(0) (date, cnt) } .reduceByKey(_ _)逻辑说明第一层reduceByKey把倾斜key拆成最多10份并行计算消除单个key的压力。第二层reduceByKey把同一日期的所有部分合并回真实结果。参数说明Random.nextInt(10)的10表示盐的基数一般设为executor总数的2倍到3倍。如果设得太大第二次聚合的shuffle量会增加反而变慢。MySQL连接数超限的原因通常是直接在rdd.foreach里DriverManager.getConnection每行建一个连接。改成foreachPartition后每个partition一个连接即可。如果executor数量很大比如200个partition也要保证MySQL的max_connections足够否则还是会被拒绝。5. 把这份代码改造成UV和转化漏斗的关键动作5.1 用distinct实现UV而不是简单计数PV口径是累加UV必须去重。源码里直接对用户ID做distinct再计数是最直接的改造点。val uvRDD rdd .map(line ((line(0), line(1)), 1)) .distinct() .map { case ((date, userId), _) (date, 1) } .reduceByKey(_ _)逻辑说明先按日期用户ID做distinct把同一个人同一天的多次访问压缩成一条再按日期计数。参数说明如果业务上要计算周UV或月UV把line(0)换成周或月的计算值即可。这里distinct产生的shuffle比较大数据量高时可以用rdd.map(line (date _ userId, 1)).reduceByKey((a, b) a).map代替语义相同但能在shuffle阶段提前去重。5.2 点击到下单再到支付的漏斗统计电商看转化漏斗通常需要统计同一session内点击某商品后是否下单、是否支付。RDD做法是先按session分组再检查动作序列。val sessionActions rdd .map(line (line(2), (line(4), line(5)))) .groupByKey() val funnelRDD sessionActions.map { case (sessionId, actions) val actionSet actions.map(_._1).toSet val clicked if (actionSet.contains(0)) 1 else 0 val ordered if (actionSet.contains(1)) 1 else 0 val paid if (actionSet.contains(2)) 1 else 0 (click, clicked) :: (order, ordered) :: (pay, paid) :: Nil }.flatMap(identity).reduceByKey(_ _)逻辑说明按session_id分组后一个session的多次动作被收集到一起。然后用集合检查动作类型是否存在避免了同一个session里多次点击导致重复计数。参数说明line(2)是session字段下标line(4)是动作类型下标line(5)是商品字段下标实际以数据格式为准。groupByKey在这一步是合理的因为每个session的数据量有限不会造成严重倾斜。5.3 把结果写成CSV而不是直接覆盖MySQL调试阶段建议把结果写成CSV格式清爽且不占用数据库连接。用repartition(1)配合saveAsTextFile输出文件再重命名。hdfs dfs -cat /output/funnel/part-00000 | head -20逻辑说明直接查看输出文件的前20行验证漏斗数字是否符合业务直觉。参数说明如果点击人数小于下单人数说明动作类型枚举值理解反了需要回到数据格式说明重新核对。最后再把确认过的结果写入MySQL避免调试期把线上统计表刷脏。本文还有配套的精品资源点击获取