如何在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
相关产品推荐
相关产品推荐

