You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在IntelliJ中用Scala 2.12.12通过Spark读取本地CSV并按月统计

解决Spark读取CSV文件及按月统计的问题

一、核心错误原因

  1. SparkSession未初始化:直接使用spark变量但未创建实例,导致上游依赖失效,df自然无法被正确定义。
  2. 路径格式错误:Windows系统下路径反斜杠未转义,会被解析为转义字符引发路径识别失败。
  3. 变量赋值逻辑错误:第二段代码中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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 14:20:23