【问题标题】:How to process files parallely in spark using spark.read function [duplicate]如何使用 spark.read 函数在 spark 中并行处理文件
【发布时间】:2018-05-24 22:09:31
【问题描述】:

我有一个包含文件列表的文本文件。目前,我正在按顺序遍历我的文件列表

我的文件列表如下所示,

D:\Users\bramasam\Documents\sampleFile1.txt
D:\Users\Documents\sampleFile2.txt

并为每个文件执行以下代码,

val df = spark.read
   .format("org.apache.spark.csv")
   .option("header", false)
   .option("inferSchema", false)
   .option("delimiter", "|")
   .schema(StructType(fields)) //calling a method to find schema
   .csv(fileName_fetched_foreach)
   .toDF(old_column_string: _*)

 df.write.format("orc").save(target_file_location)

我要做的是为每个文件而不是序列并行执行上述代码,因为文件之间没有依赖关系。所以,我正在尝试类似下面的方法,但遇到了错误,

  //read the file which has the file list
    spark.read.textFile("D:\\Users\\Documents\\ORC\\fileList.txt").foreach { line =>
      val tempTableName = line.substring(line.lastIndexOf("\\"),line.lastIndexOf("."))
      val df = spark.read
        .format("org.apache.spark.csv")
        .option("header", false)
        .option("inferSchema", false)
        .option("delimiter", "|")
        .schema(StructType(fields))
        .csv(line)
        .toDF(old_column_string: _*)
        .registerTempTable(tempTableName)

      val result = spark.sql(s"select $new_column_string from $tempTableName") //reordering column order on how it has to be stored

      //Note: writing to ORC needs Hive support. So, make sure the systax is right
      result.write.format("orc").save("D:\\Users\\bramasam\\Documents\\SCB\\ORCFile")
    }
  }

我面临以下错误,

java.lang.NullPointerException
    at org.apache.spark.sql.SparkSession.sessionState$lzycompute(SparkSession.scala:135)
    at org.apache.spark.sql.SparkSession.sessionState(SparkSession.scala:133)
    at org.apache.spark.sql.DataFrameReader.<init>(DataFrameReader.scala:689)
    at org.apache.spark.sql.SparkSession.read(SparkSession.scala:645)
    at ConvertToOrc$$anonfun$main$1.apply(ConvertToOrc.scala:25)
    at ConvertToOrc$$anonfun$main$1.apply(ConvertToOrc.scala:23)
    at scala.collection.Iterator$class.foreach(Iterator.scala:727)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1157)

【问题讨论】:

  • 您能否详细说明您正在尝试做什么以及出于什么目的?
  • @ShankarKoirala,我正在尝试处理一组文件(在这种情况下,通过应用架构将文件转换为 ORC)。但是,在我当前的过程中,我通过运行 for 循环遍历文件名来顺序读取文件。由于过程是独立的,我想通过使用 spark 的高阶函数来并行实现这一点。

标签: scala apache-spark foreach apache-spark-sql


【解决方案1】:

你应该做的是将文件放在一个目录中,然后让 spark 读取整个目录,如果需要,在每个文件中附加一列并带有文件名

spark.read.textFile("D:\\Users\\Documents\\ORC\\*")

    • 是读取所有文件

【讨论】:

    【解决方案2】:

    请查看@samthebest answer

    在您的情况下,您应该传递整个目录,例如:

    spark.read.textFile("D:\\Users\\Documents\\ORC")
    

    如果您想递归读取目录,请参阅 to this answer:

    spark.read.textFile("D:\\Users\\Documents\\ORC\\*\\*")
    

    【讨论】:

      猜你喜欢
      • 2016-07-31
      • 2016-05-26
      • 2018-11-09
      • 2015-09-05
      • 1970-01-01
      • 1970-01-01
      • 2022-07-05
      • 1970-01-01
      • 2018-08-16
      相关资源
      最近更新 更多