基于Spark RDD实现电商用户购买行为数据分析
前言
Hello,我又来更新啦,截止到现在,有关于Hadoop和spark相关的安装部署就基本上结束啦,那么今天以基于Spark RDD实现电商用户购买行为数据分析为项目,带着大家走一遍spark的应用,同样每一步都有详细的文字介绍以及图片展示,清晰易懂,如果你也感兴趣这篇文章的话,那就跟随我的脚步往下看吧~
同时如果这篇文章对你有所帮助,或者你也喜欢的话,可以动动你发财的小手给作者点点赞+收藏吗❤️你的一个赞就是我持续学习更新的动力~
目录
项目任务和内容
某电商平台积累了大量原始用户购买日志数据,这些数据直接从业务系统采集,存在格式不统一、无效值、重复记录等问题,无法直接用于业务分析。现需要使用 Spark RDD 技术对该批数据进行预处理与多维度聚合,提取有效业务指标,为后续商品运营、用户画像分析提供数据支撑。
基于 Spark RDD 实现电商用户购买行为数据的全流程清洗、标准化与多维度聚合分析。
1.数据来源
本地文本文件 user_purchase_logs.txt,每行对应一条原始用户购买日志,文件编码为 UTF-8。
2.原始数据格式详情
原始 RDD 加载后,每条数据为一个字符串(每行对应一个元素),字段之间使用英文逗号分隔,格式定义为:用户ID,商品ID,购买时间字符串,购买数量,支付金额字符串各字段详细说明:
| 字段名称 | 字段含义 | 数据特征与可能取值 |
| 用户 ID | 唯一标识电商平台用户 | 格式为user_xxx,其中xxx为数字(如user_001、user_123) |
| 商品 ID | 唯一标识平台商品 | 格式为goods_xxx/elect_xxx/cloth_xxx,前缀为商品类别(如goods_005、elect_010) |
| 购买时间字符串 | 用户完成购买的时间记录 | 格式不统一,包含两种:① 仅日期(yyyy-MM-dd);② 日期 + 时分(yyyy/MM/dd HH:mm) |
| 购买数量 | 单次订单中该商品的购买件数 | 可能为正整数、负整数(无效订单)、零(无效订单) |
| 支付金额字符串 | 单次订单的支付总金额 | 可能为带单位字符串(xxx.yy元)、纯数字字符串(xxx.yy)、无效字符串(无效金额、待支付) |
3.原始样本数据(可自行添加)
| user_001,goods_001,2025-12-25,2,199.99元 user_001,goods_001,2025-12-25,2,199.99元 user_002,elect_005,2025/12/26 09:30,-1,5999.00元 user_003,cloth_010,2025-12-27,5,299.99 user_004,goods_002,2025/12/25 14:20,0,无效金额 user_005,elect_006,2025-12-26,3,699.00元 user_002,elect_005,2025/12/26 10:00,1,5999.00元 user_006,cloth_011,2025/12/27,8,待支付 user_003,goods_003,2025/12/27 16:40,2,89.99元 user_007,elect_007,2025-12-24,4,1299.00 user_001,cloth_010,2025/12/25 18:15,3,399.99元 user_008,goods_004,2025-12-26,-5,39.99元 user_009,elect_005,2025/12/27,1,5999.00元 user_005,cloth_012,2025/12/26 11:20,2,259.99 user_003,goods_001,2025-12-27,1,199.99元 user_007,goods_002,2025-12-24 09:00,5,79.98元 user_002,cloth_010,2025/12/26,2,399.99元 user_004,elect_008,2025/12/25,3,1999.00元 user_006,goods_003,2025/12/27,2,89.99元 user_009,cloth_011,2025/12/27 19:00,4,499.99元 |
一、前置操作
1. 启动 Hadoop
这一步是必须先启动的,如果Hadoop安装没有完善的可以移步我的前几篇文章去查看安装步骤。
进入hadoop目录
cd /usr/local/hadoop
启动HDFS+YARN
start-all.sh
验证进程
jps
如下图,正常输出需包含:NameNode、DataNode、ResourceManager、NodeManager,全部出现说明你的Hadoop部署没有问题。



2.启动 Spark 集群
这一步如果spark的相关安装部署没有完善的话,可以移步我的前几篇文章去查看详细步骤学习。
进入Spark根目录
cd /export/servers/spark
一键启动Standalone集群
sbin/start-all.sh

验证Spark进程
jps
如下图,hadoop1 正常进程出现Master、Worker,hadoop2/hadoop3 正常进程出现Worker,代表没有问题。



3.创建数据文件夹并写入数据集
执行这条命令,一键创建目录并打开文件:
mkdir -p /usr/local/data
vim /usr/local/data/user_purchase_logs.txt
按 i 进入编辑模式,输入原始全部数据。
输入完成后按 Esc,输入 :wq 保存退出。


4.启动 spark-shell
执行绝对路径命令:
/export/servers/spark/bin/spark-shell --master local[2]
如下图,出现 scala> 代表进入成功。

二、Spark Shell 完整分步代码
1.加载原始数据
输入读取的文件路径如下:
val logPath = "file:///usr/local/data/user_purchase_logs.txt"
val rawPurchaseRDD = sc.textFile(logPath)
println("原始日志总条数:" + rawPurchaseRDD.count() )
rawPurchaseRDD.take(5).foreach(println)

2.脏数据过滤
这一步是为了过滤无效的订单。
val validPurchaseRDD = rawPurchaseRDD.filter { line =>
val fields = line.split(",")
if (fields.length != 5) false
else {
val numStr = fields(3).trim
val moneyStr = fields(4).trim
var isValid = true
try {
val buyNum = numStr.toInt
if (buyNum <= 0) isValid = false
val cleanMoney = moneyStr.replace("元", "")
if (cleanMoney == "无效金额" || cleanMoney == "待支付") isValid = false
else cleanMoney.toDouble
} catch {
case _: Exception => isValid = false
}
isValid
}
}

println("\n过滤后有效数据条数:" + validPurchaseRDD.count() )
validPurchaseRDD.collect().foreach(println)

3.数据标准化
这一步是为了统一时间、数值格式。
val standardizedPurchaseRDD = validPurchaseRDD.map { line =>
val fields = line.split(",")
val userId = fields(0).trim
val goodsId = fields(1).trim
val timeStr = fields(2).trim
val buyNum = fields(3).trim.toInt
val buyMoney = fields(4).trim.replace("元", "").toDouble
// 统一时间为 yyyy-MM-dd HH:mm:ss
val standardTime = if (timeStr.contains("/")) {
timeStr.replace("/", "-") + ":00"
} else {
timeStr + " 00:00:00"
}
(userId, goodsId, standardTime, buyNum, buyMoney)
}
println("标准化样例数据")
standardizedPurchaseRDD.take(10).foreach(println)

4.重复订单去重
这一步是利用用户 + 商品 + 时间唯一键去重。
val uniquePurchaseRDD = standardizedPurchaseRDD
.map(t => (t._1 + "|" + t._2 + "|" + t._3, t))
.distinct()
.map(_._2)
println("去重前:" + standardizedPurchaseRDD.count() + " 条,去重后:" + uniquePurchaseRDD.count() + " 条")

5.维度扩展,提取商品大类
val extendedPurchaseRDD = uniquePurchaseRDD.map { t =>
val goodsType = t._2.split("_")(0)
// key:(用户ID,商品类别) value:(商品ID,购买数量,金额,订单计数1)
((t._1, goodsType), (t._2, t._4, t._5, 1))
}
println("维度扩展后数据")
extendedPurchaseRDD.collect().foreach(println)

6.分组聚合统计
val aggregatedPurchaseRDD = extendedPurchaseRDD
.reduceByKey { (v1, v2) =>
val totalNum = v1._2 + v2._2
val totalMoney = v1._3 + v2._3
val orderCnt = v1._4 + v2._4
(v1._1, totalNum, totalMoney, orderCnt)
}
.map { case ((uid, gtype), (gid, totalNum, totalMoney, orderCnt)) =>
val avgNum = (totalNum.toDouble / orderCnt).formatted("%.2f").toDouble
// 输出结构:(用户ID,商品类别,总金额,总件数,平均单次购买量)
(uid, gtype, totalMoney, totalNum, avgNum)
}
println("聚合统计结果")
aggregatedPurchaseRDD.collect().foreach(println)

7.结果保存到本地文件
输出路径为/usr/local/data下,代码如下:
val outputPath = "file:///usr/local/data/purchase_result"

为了防止报错,先删除旧结果文件夹。
import org.apache.hadoop.fs.{FileSystem, Path}
val fs = FileSystem.get(sc.hadoopConfiguration)
if(fs.exists(new Path(outputPath))){
fs.delete(new Path(outputPath),true)
}

格式化输出并保存。
aggregatedPurchaseRDD
.map{ case ((uid, gt), (gid, num, money, orderCnt)) =>
// 计算平均值并保留2位小数
val avg = (num.toDouble / orderCnt).formatted("%.2f")
s"用户:${uid},商品类别:${gt},总消费:${money}元,总件数:${num},平均单次购买:${avg}件"
}
.coalesce(1)
.saveAsTextFile(outputPath)

打印最终报告。
aggregatedPurchaseRDD.foreach{ case ((uid, gt), (gid, num, money, orderCnt)) =>
// 计算平均购买数量,保留2位小数
val avg = (num.toDouble / orderCnt).formatted("%.2f")
println(s"""
用户ID:$uid
商品分类:$gt
累计消费金额:$money 元
累计购买件数:$num 件
平均单次购买量:$avg 件
----------------------------------------
""")
}



三、实验后查看结果
退出 spark-shell 后执行查看命令。
进入结果目录
cd /usr/local/data/purchase_result
查看结果文件
cat part-00000

总结
在应用spark做实验或者项目的过程中,大家如果也有遇到的问题,可以留言评论,我们互相学习~
以上就是个人对于在实验项目中应用spark的每一步操作,希望可以帮助看到这篇文章的你,如有不足,希望大家多多包涵,有错的地方也欢迎大家指出,多多指教,我们互相学习,仅以此篇,与君共勉。
我是正在持续学习进步中的小丁同学哇~希望大家可以给我个点赞+收藏。
更多推荐
所有评论(0)