【问题标题】:Formatting the join rdd - Apache Spark格式化连接 rdd - Apache Spark
【发布时间】:2015-06-28 01:03:02
【问题描述】:

我有两个键值对 RDD,我加入两个 rdd 并保存为文本文件,代码如下:

val enKeyValuePair1 = rows_filter6.map(line => (line(8) -> (line(0),line(4),line(10),line(5),line(6),line(14),line(1),line(9),line(12),line(13),line(3),line(15),line(7),line(16),line(2),line(14))))

val enKeyValuePair = DATA.map(line => (line(0) -> (line(2),line(3))))

val final_res = enKeyValuePair1.leftOuterJoin(enKeyValuePair)

val output = final_res.saveAsTextFile("C:/out")

my output is as follows:
(534309,((17999,5161,45005,00000,XYZ,,29.95,0.00),None))

我怎样才能去掉所有的括号? 我希望我的输出如下:

534309,17999,5161,45005,00000,XYZ,,29.95,0.00,None

【问题讨论】:

    标签: scala join apache-spark rdd keyvaluepair


    【解决方案1】:

    当输出到文本文件时,Spark 将只使用 RDD 中元素的toString 表示。如果您想控制格式,则可以在调用saveAsTextFile 之前将数据最后一次转换为String

    幸运的是,使用 Spark API 生成的元组可以使用解构来拆分。在您的示例中,我会这样做:

    val final_res = enKeyValuePair1.leftOuterJoin(enKeyValuePair)
    val formatted = final_res.map { tuple =>
      val (f1,((f2,f3,f4,f5,f6,f7,f8,f9),f10)) = tuple
      Seq(f1, f2, f3, f4, f5, f6, f7, f8, f9, f10).mkString(",")
    }
    formatted.saveAsTextFile("C:/out")
    

    第一行val 将采用传递给map 函数的元组并将组件分配给左侧的值。第二行创建一个临时的Seq,其中包含您希望显示的顺序的字段,然后调用mkString(",") 以使用逗号连接字段。

    如果字段较少或您只是解决 REPL 上的问题,也可以通过对传递给 map 的部分函数使用模式匹配来稍微替代上述方法。

    simpleJoinedRdd.map { case (key,(left,right)) => s"$key,$left,$right"}}
    

    虽然这确实允许您将其设为单行表达式,但如果 R​​DD 中的数据与提供的模式不匹配,它可能会引发异常,这与前面的示例相反,如果 tuple 参数编译器将抱怨无法解构为预期的形式。

    【讨论】:

      【解决方案2】:

      你可以这样做:

      import scala.collection.JavaConversions._
      val output = sc.parallelize(List((534309,((17999,5161,45005,1,"XYZ","",29.95,0.00),None))))
      val result = output.map(p => p._1 +=: p._2._1.productIterator.toBuffer += p._2._2)
        .map(p => com.google.common.base.Joiner.on(", ").join(p.iterator))
      

      我使用番石榴来格式化字符串,但有可能是 scala 的方式。

      【讨论】:

        【解决方案3】:

        保存前做一个平面图。或者,您可以编写一个简单的格式化函数并在 map 中使用它。 添加一点代码,只是为了展示它是如何完成的。函数 formatOnDemand 可以是任何东西

        test = sc.parallelize([(534309,((17999,5161,45005,00000,"XYZ","",29.95,0.00),None))])
        print test.collect()
        print test.map(formatOnDemand).collect()
        
        def formatOnDemand(t):
            out=[]
            out.append(t[0])
            for tok in t[1][0]:
                out.append(tok)
            out.append(t[1][1])
            return out
        
        >>> 
        [(534309, ((17999, 5161, 45005, 0, 'XYZ', '', 29.95, 0.0), None))]
        [[534309, 17999, 5161, 45005, 0, 'XYZ', '', 29.95, 0.0, None]]
        

        【讨论】:

        • 你能提供更多细节吗?
        猜你喜欢
        • 1970-01-01
        • 2014-05-13
        • 1970-01-01
        • 2017-03-07
        • 1970-01-01
        • 1970-01-01
        • 2018-04-08
        • 2018-03-29
        • 2014-11-16
        相关资源
        最近更新 更多