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

如何在Spark中用spark.read并行处理文件及解决报错问题

问题分析与解决方案

首先得说清楚你为什么会遇到这个NullPointerException:Spark的foreach是在Executor节点的Task中执行的,而你创建的SparkSession实例(也就是代码里的spark对象)是在Driver节点上的,它并没有被序列化传递到Executor端。所以当Executor里的Task尝试调用spark.read时,spark对象是null,自然就抛出空指针异常了。

而且你其实走进了一个误区:Spark本身就是为并行处理设计的,根本不需要你手动用foreach去实现并行——Spark读取文件时会自动并行处理多个文件,我们只要利用它的原生能力就好。下面给你两种可行的解决方案:

方案一:直接批量读取所有文件(适合文件数量不多的场景)

这种方式最简单,直接把所有文件路径传给csv方法,Spark会自动并行读取这些文件,完全不需要手动循环:

// 读取文件列表,提取所有文件路径
val filePaths = spark.read.textFile("D:\\Users\\Documents\\ORC\\fileList.txt")
  .collect() // 注意:如果文件列表特别大,不要用collect,会导致Driver内存溢出,改用方案二

// 一次性读取所有文件,Spark自动并行处理
val rawDF = spark.read
  .format("org.apache.spark.csv")
  .option("header", false)
  .option("inferSchema", false)
  .option("delimiter", "|")
  .schema(StructType(fields))
  .csv(filePaths: _*)
  .toDF(old_column_string: _*)

// 重新调整列顺序
val resultDF = rawDF.selectExpr(new_column_string.split(","): _*)

// 写入ORC格式
resultDF.write.format("orc").save("D:\\Users\\bramasam\\Documents\\SCB\\ORCFile")

方案二:分区处理大文件列表(适合文件数量极多的场景)

如果你的文件列表非常大,collect()会把所有路径加载到Driver内存里导致OOM,这时候可以用mapPartitions在Executor端的分区里批量处理文件:

import org.apache.spark.sql.SparkSession

spark.read.textFile("D:\\Users\\Documents\\ORC\\fileList.txt")
  .mapPartitions { filePathIter =>
    // 在每个分区的Task中创建/获取SparkSession(Spark 2.x及以上支持这种方式)
    val localSpark = SparkSession.builder()
      .config(spark.sparkContext.getConf)
      .getOrCreate()
    import localSpark.implicits._

    // 遍历当前分区的所有文件路径,逐个读取并处理
    filePathIter.flatMap { path =>
      val df = localSpark.read
        .format("org.apache.spark.csv")
        .option("header", false)
        .option("inferSchema", false)
        .option("delimiter", "|")
        .schema(StructType(fields))
        .csv(path)
        .toDF(old_column_string: _*)
        .selectExpr(new_column_string.split(","): _*)
      // 将DataFrame转为本地迭代器,方便合并分区结果
      df.toLocalIterator()
    }
  }
  .toDF()
  .write.format("orc").save("D:\\Users\\bramasam\\Documents\\SCB\\ORCFile")

额外注意事项

  • 永远不要在foreach、map这类Executor端执行的算子里直接引用Driver端的SparkSession、SparkContext对象,它们无法被序列化传输到Executor,必然会导致空指针或者序列化错误。
  • 如果需要每个文件生成单独的ORC文件,可以在处理时给每个文件的DataFrame加上一个标识列,然后用partitionBy写入,或者在mapPartitions里单独写入每个文件到不同路径。

内容的提问来源于stack exchange,提问作者Gladiator

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 10:08:06