在IntelliJ中用Scala 2.12.12通过Spark读取本地CSV并按月统计
解决Spark读取CSV文件及按月统计的问题
一、核心错误原因
- SparkSession未初始化:直接使用
spark变量但未创建实例,导致上游依赖失效,df自然无法被正确定义。 - 路径格式错误:Windows系统下路径反斜杠未转义,会被解析为转义字符引发路径识别失败。
- 变量赋值逻辑错误:第二段代码中
df未提前声明为可变变量,且.show()返回Unit类型,无法赋值给DataFrame变量。
二、正确实现方案
1. 适配Scala 2.12的build.sbt依赖配置
name := "SparkCSVProcessor" version := "0.1" scalaVersion := "2.12.12" libraryDependencies ++= Seq( "org.apache.spark" %% "spark-sql" % "3.3.0" // 3.3.x版本完美适配Scala 2.12 )
2. 完整Scala代码(含读取与统计逻辑)
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ object CSVDataProcessor { def main(args: Array[String]): Unit = { // 初始化SparkSession核心实例 val spark = SparkSession.builder() .appName("CSVMonthlyStats") .master("local[*]") // 本地调试启用,生产环境移除该行 .getOrCreate() // 导入隐式转换,简化DataFrame操作 import spark.implicits._ // 处理Windows路径:用双反斜杠转义,或改用正斜杠 val csvPath = "C:\\Users\\trialrun\\Desktop\\DataExtract.csv" // 替代写法:val csvPath = "C:/Users/trialrun/Desktop/DataExtract.csv" // 读取带表头的CSV文件,自动推断列类型 val df = spark.read .option("header", "true") .option("inferSchema", "true") // 必须开启,否则数值列会被识别为字符串,无法求和 .csv(csvPath) // 验证读取结果 df.show(5) // 按月统计总和(假设日期列名为`date`,待求和列名为`amount`,请根据实际列名修改) val monthlySumResult = df .withColumn("month", date_format(col("date"), "yyyy-MM")) .groupBy("month") .agg(sum("amount").alias("total_amount")) // 输出统计结果 monthlySumResult.show() // 关闭SparkSession spark.stop() } }
三、关键注意事项
- SparkSession是核心:所有Spark操作必须基于该实例,本地调试用
local[*]可利用多线程加速百万行数据处理。 - 路径转义:Windows路径必须用双反斜杠或正斜杠,避免解析错误。
- Schema推断:开启
inferSchema才能让Spark识别数值类型,否则求和操作会报错。 - 变量声明:用
val定义DataFrame(不可变特性),不要用var随意赋值,避免类型冲突。
内容的提问来源于stack exchange,提问作者Kacie
相关产品推荐
相关产品推荐

