【问题标题】:Format (remove class/parens) Spark CSV saveAsTextFile output?格式化(删除类/括号)Spark CSV saveAsTextFile 输出?
【发布时间】:2015-07-16 21:09:10
【问题描述】:

我正在尝试从通过 saveAsTextFile 保存的 CSV 数据中剥离包装类或数组文本,而无需执行非 Spark 后处理步骤。

我有一些大文件中的 TSV 数据,我将其提供给 RDD。

 val testRdd = sc.textFile(_input).filter(!_.startsWith("unique_transaction_id")).map(x => x.toLowerCase).map(x => x.split('\t')).map(x => Test(x(0),x(1)))

testRdd.saveAsTextFile("test")

这样保存了类名包裹的数据:

head -n 1 part-00000
Test("1969720fb3100608b38297aad8b3be93","active")

我还尝试将其用于未命名的类 (?) 而不是案例类。

val testRdd = sc.textFile(_input).filter(!_.startsWith("unique_transaction_id")).map(x => x.toLowerCase).map(x => x.split('\t')).map(x => (x(0),x(1)))

testRdd.saveAsTextFile("test2")

这会产生

("1969720fb3100608b38297aad8b3be93","active")

仍然需要后处理才能删除包装括号。

为了去除包装字符,我尝试了 flatMap(),但 RDD 显然不是正确的类型:

testRdd.flatMap(identity).saveAsTextFile("test3")
<console>:17: error: type mismatch;
 found   : ((String, String)) => (String, String)
 required: ((String, String)) => TraversableOnce[?]
              testRdd.flatMap(identity).saveAsTextFile("test3")

那么...我需要将 RDD 转换为其他类型的 RDD,还是有另一种方法可以将 RDD 保存为 CSV 以便剥离换行文本?

谢谢!

【问题讨论】:

    标签: csv apache-spark rdd


    【解决方案1】:
    val testRdd = sc.textFile(_input).filter(!_.startsWith("unique_transaction_id")).map(x => x.toLowerCase).map(x => x.split('\t')).map(x => x(0)+","+x(1))
    

    这会将输出写为 csv

    【讨论】:

      【解决方案2】:

      您可以尝试以下方法:

      val testRdd = sc.textFile(_input).filter(!_.startsWith("unique_transaction_id"))
                                       .map(x => x.toLowerCase.split('\t'))
                                       .map(x => x(0)+","+x(1))
      

      我们听到的是过滤你的标题后,你可以在同一个地图段落中小写你的字符串,还可以节省一些不必要的额外映射。

      这将创建一个可以保存为 CSV 格式的 RDD[String]。

      PS:

      • 保存的rdd输出的扩展名不是csv而是格式!

      • 这不是最佳且唯一的解决方案,但它会为您完成这项工作!

      【讨论】:

      • 完美。感谢您不仅回答了这个问题,还指出了在单个 map() 调用中执行多个步骤的简化。
      【解决方案3】:

      你可以看看Spark CSV Library

      【讨论】:

        【解决方案4】:

        val logFile = "/input.csv"

        val conf = new SparkConf().set("spark.driver.allowMultipleContexts", "true")

        val sc = new SparkContext(master="local", appName="Mi app", conf)

        val logData = sc.textFile(logFile, 2).cache()

        val lower = logData.map(line => line.toLowerCase)

        【讨论】:

          猜你喜欢
          • 2016-09-18
          • 1970-01-01
          • 2023-01-15
          • 2019-06-10
          • 2016-09-18
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2012-11-25
          相关资源
          最近更新 更多