前言

Hello,我又来更新啦,截止到现在,有关于Hadoop和spark相关的安装部署就基本上结束啦,那么今天以基于Spark RDD实现电商用户购买行为数据分析为项目,带着大家走一遍spark的应用,同样每一步都有详细的文字介绍以及图片展示,清晰易懂,如果你也感兴趣这篇文章的话,那就跟随我的脚步往下看吧~

同时如果这篇文章对你有所帮助,或者你也喜欢的话,可以动动你发财的小手给作者点点赞+收藏吗❤️你的一个赞就是我持续学习更新的动力~
 

目录

前言

项目任务和内容

一、前置操作

1. 启动 Hadoop

2.启动 Spark 集群

3.创建数据文件夹并写入数据集

4.启动 spark-shell

二、Spark Shell 完整分步代码

1.加载原始数据

2.脏数据过滤

3.数据标准化

4.重复订单去重

5.维度扩展,提取商品大类

6.分组聚合统计

7.结果保存到本地文件

三、实验后查看结果

总结


项目任务和内容

某电商平台积累了大量原始用户购买日志数据,这些数据直接从业务系统采集,存在格式不统一、无效值、重复记录等问题,无法直接用于业务分析。现需要使用 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的每一步操作,希望可以帮助看到这篇文章的你,如有不足,希望大家多多包涵,有错的地方也欢迎大家指出,多多指教,我们互相学习,仅以此篇,与君共勉。

我是正在持续学习进步中的小丁同学哇~希望大家可以给我个点赞+收藏。

更多推荐