【问题标题】:Hadoop Spark 1.4.1 - sort multiple CSV files and save sorted result in 1 output fileHadoop Spark 1.4.1 - 对多个 CSV 文件进行排序并将排序结果保存在 1 个输出文件中
【发布时间】:2016-06-27 04:05:55
【问题描述】:

我在 HDFS 中有 3 个文件,我想使用最有效的方法先在第一列排序,然后在第二列排序,然后在 Spark 1.4 中使用 Scala(或 Python)将排序结果存储回 HDFS 上的新文件。 1:
hdfs:///test/2016/file.csv
hdfs:///test/2015/file.csv
hdfs:///test/2014/file.csv

文件看起来像这样(没有标题):
hdfs:///test/2016/file.csv
127,56,abc
125,56,abc
121,56,abc

hdfs:///test/2016/file.csv
126,66,abc
122,56,abc
123,46,abc

hdfs:///test/2016/file.csv
122,66,abc
128,56,abc
123,16,abc

排序后的输出要保存到 HDFS:
hdfs:///test/output/file.csv
121,56,abc
122,56,abc
122,66,abc
123,16,abc
123,46,abc
125,56,abc
126,66,abc
127,56,abc
128,56,abc

我对 Spark 很陌生,到目前为止我只知道如何加载文件:
val textFile = sc.textFile("hdfs:///test/2016/file.csv")

试图在互联网上阅读如何排序,但不清楚哪些库应该适用于这种情况(CSV 文件)和这个版本的 Spark(1.4.1)以及如何使用它们。 请帮忙,乔

【问题讨论】:

    标签: python scala csv hadoop apache-spark


    【解决方案1】:

    我建议使用 databricks csv 库来读写 csv:https://github.com/databricks/spark-csv

    由于我现在无法访问 hdfs,因此此示例使用文件系统,但在与 hdfs 路径一起使用时也应该可以工作。

    import org.apache.spark.sql.SQLContext
    import org.apache.spark.{SparkContext, SparkConf}
    import org.apache.spark.sql.functions._  // needed for ordering the dataframe
    
    object StackoverflowTest {
    
      def main(args: Array[String]) {
        // basic spark init
        val conf = new SparkConf().setAppName("Data Import from CSV")
        val sc = new SparkContext(conf)
        val sqlContext = new SQLContext(sc)
    
        // first we load every file from the data directory that starts with 'file'
        val storeDf = sqlContext.read
          .format("com.databricks.spark.csv")
          .option("inferSchema", "true")
          .load("data/file*")
    
        // then we sort it and write to an output
        storeDf
          .orderBy("C0", "C1")  // default column names
          .repartition(1)   // in order to have 1 output file
          .write
          .format("com.databricks.spark.csv")
          .save("data/output")
      }
    
    }
    

    结果将作为 csv 写入 data/output/part-00000。 希望这会有所帮助。

    【讨论】:

    • 这适用于 Spark 1.4.1 版吗?我正在使用控制台,在这一行之后我看到 Error: val sc = new SparkContext(conf) org.apache.spark.SparkException: 只有一个 SparkContext 可能在这个 JVM 中运行(参见 SPARK-2243)。要忽略此错误,请设置 spark.driver.allowMultipleContexts = true。当前运行的 SparkContext 创建于:org.apache.spark.SparkContext.(SparkContext.scala:81) org.apache.spark.repl.SparkILoop.createSparkContext(SparkILoop.scala:1017)
    • 那是因为您已经在该 shell 中创建了 SparkContext - 要么关闭此 shell 会话并启动另一个会话,要么跳过上下文创建并使用您已经创建的会话。
    • 是的,正如@TzachZohar 已经指出的那样,如果你在 shell 中运行它,你不需要 val conf、val sc、val sqlContext 行,它们已经默认初始化了
    • 在执行 .load("correct hdfs path") 行后 - 这是错误:java.lang.RuntimeException:无法加载数据源的类:scala.sys 的 com.databricks.spark.csv .package$.error(package.scala:27​​) at org.apache.spark.sql.sources.ResolvedDataSource$.lookupDataSource(ddl.scala:220) at org.apache.spark.sql.sources.ResolvedDataSource$.apply( ddl.scala:233) 在 org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:114)
    • 那是因为它需要 databricks csv 库,正如我发布的链接提到的那样,如果您将它与 spark shell 一起使用,则需要使用 spark-shell --packages com 启动 shell。 databricks:spark-csv_2.11:1.4.0 或 spark-shell --packages com.databricks:spark-csv_2.10:1.4.0 取决于您使用的是 scala 2.10 还是 2.11
    【解决方案2】:
    val textFile = sc.textFile("hdfs:///test/*/*.csv")
                     .map( _.split(",",-1) match { case Array(col1, col2, col3) => (col1, col2, col3) })
                     .sortBy(_._1)
                     .map(_._1+","+_._2+","+_._3)
                     .saveAsTextFile("hdfs:///testoutput/output/file.csv")
    

    您需要保存在不同的文件夹中,否则您生成的文件将在您再次运行时被重复使用。

    【讨论】:

    • 感谢分享。我正在使用控制台并尝试此操作。在尝试此方法之前是否需要导入任何库?您如何测试排序阶段是否良好? (对我来说,它在第 3 行之后显示空 - 没有结果 - 当我这样做时 sortBy(_._1):textFile.foreach(println)
    • @Joe 我有一个错字。应该是sc.textFile("hdfs:///test/*/*.csv")。你纠正了吗?您不需要任何库
    • 我看到了小字体并已修复,但它仍然对我不起作用。这是来自我的控制台:scala> .map(.1+","+。 _2+","+._3) :1: error: ')' 预期但发现双重文字。 res9.map(.1+","+._2+","+._3)
    • res9.map(_._1+","+_._2+","+_._3).
    • 你能运行这个val textFile = sc.textFile("hdfs:///test/*/*.csv").count() 并告诉我结果吗?
    猜你喜欢
    • 2017-08-01
    • 2017-07-17
    • 1970-01-01
    • 1970-01-01
    • 2014-03-30
    • 2011-12-03
    • 2021-09-08
    • 1970-01-01
    • 2017-06-23
    相关资源
    最近更新 更多