【问题标题】:Create Tuple out of Array(Array[String) of Varying Sizes using Scala使用 Scala 从不同大小的数组(数组 [字符串)中创建元组
【发布时间】:2019-04-22 10:16:08
【问题描述】:

我是 scala 的新手,我正在尝试用一个 Array(Array[String]) 类型的 RDD 来创建一个 Tuple 对,它看起来像:

(122abc,223cde,334vbn,445das),(221bca,321dsa),(231dsa,653asd,698poq,897qwa)

我正在尝试从这些数组中创建元组对,以便每个数组的第一个元素是键,而数组的任何其他部分都是一个值。例如,输出如下所示:

122abc    223cde
122abc    334vbn
122abc    445das
221bca    321dsa
231dsa    653asd
231dsa    698poq
231dsa    897qwa

我不知道如何从每个数组中分离第一个元素,然后将其映射到其他所有元素。

【问题讨论】:

  • 为什么你有两个221bca 321dsa
  • @smac89 这是一个错字对不起。现在改了。
  • 您是否尝试将RDD[Array[Array[String]]] 映射到RDD[(String,String)]
  • @JackLeow 是的,我正在尝试将 RDD[Array[Array[String]]] 映射到 RDD[(String,String)]。对不起,如果我不够清楚。

标签: arrays scala apache-spark rdd


【解决方案1】:

如果我没看错,您的问题的核心与将内部数组的头部(第一个元素)与尾部(其余元素)分开有关,您可以使用 head 和 @987654322 @ 方法。 RDD 的行为很像 Scala 列表,因此您可以使用看起来像纯 Scala 代码的代码来完成这一切。

给定以下输入 RDD:

val input: RDD[Array[Array[String]]] = sc.parallelize(
  Seq(
    Array(
      Array("122abc","223cde","334vbn","445das"),
      Array("221bca","321dsa"),
      Array("231dsa","653asd","698poq","897qwa")
    )
  )
)

以下应该做你想做的事:

val output: RDD[(String,String)] =
  input.flatMap { arrArrStr: Array[Array[String]] =>
    arrArrStr.flatMap { arrStrs: Array[String] =>
      arrStrs.tail.map { value => arrStrs.head -> value }
    }
  }

事实上,由于flatMap/map 的构成方式,您可以将其重写为便于理解。:

val output: RDD[(String,String)] =
  for {
    arrArrStr: Array[Array[String]] <- input
    arrStr: Array[String] <- arrArrStr
    str: String <- arrStr.tail
  } yield (arrStr.head -> str)

你选择哪一个最终取决于个人喜好(尽管在这种情况下,我更喜欢后者,因为你不必缩进那么多代码)。

验证:

output.collect().foreach(println)

应该打印出来:

(122abc,223cde)
(122abc,334vbn)
(122abc,445das)
(221bca,321dsa)
(231dsa,653asd)
(231dsa,698poq)
(231dsa,897qwa)

【讨论】:

    【解决方案2】:

    这是一个经典的折叠操作;但在 Spark 中折叠调用 aggregate:

    // Start with an empty array
    data.aggregate(Array.empty[(String, String)]) { 
      // `arr.drop(1).map(e => (arr.head, e))` will create tuples of 
      // all elements in each row and the first element.
      // Append this to the aggregate array.
      case (acc, arr) => acc ++ arr.drop(1).map(e => (arr.head, e))
    }
    

    解决方案是非 Spark 环境:

    scala> val data = Array(Array("122abc","223cde","334vbn","445das"),Array("221bca","321dsa"),Array("231dsa","653asd","698poq","897qwa"))
    scala> data.foldLeft(Array.empty[(String, String)]) { case (acc, arr) =>
         |     acc ++ arr.drop(1).map(e => (arr.head, e))
         | }
    res0: Array[(String, String)] = Array((122abc,223cde), (122abc,334vbn), (122abc,445das), (221bca,321dsa), (231dsa,653asd), (231dsa,698poq), (231dsa,897qwa))
    

    【讨论】:

      【解决方案3】:

      将您的输入元素转换为 seq 和 all 然后尝试编写包装器,它将为您提供List(List(item1,item2), List(item1,item2),...)

      试试下面的代码

      val seqs = Seq("122abc","223cde","334vbn","445das")++
      Seq("221bca","321dsa")++
      Seq("231dsa","653asd","698poq","897qwa")
      

      编写一个包装器将 seq 转换为一对二

      def toPairs[A](xs: Seq[A]): Seq[(A,A)] = xs.zip(xs.tail)
      

      现在将您的 seq 作为参数发送,它将给您一对两个

      toPairs(seqs).mkString(" ")
      

      将其转换为字符串后,您将获得类似

      的输出
      res8: String = (122abc,223cde) (223cde,334vbn) (334vbn,445das) (445das,221bca) (221bca,321dsa) (321dsa,231dsa) (231dsa,653asd) (653asd,698poq) (698poq,897qwa)
      

      现在您可以随意转换字符串了。

      【讨论】:

      • 我不确定,但你的输出看起来不像 OP。
      • toPairs(seqs) 会给你List(List(item1,item2),List(item1,item2)...) 所以它应该是差不多的,然后你可以转换成你想要的。
      • 不,这不是 OP 想要的。 OP 希望创建一个元组数组,其中元组来自每个子数组的第一个元素,以及原始 RDD 中每个子数组的子数组的其余元素。
      【解决方案4】:

      使用 df 并爆炸。

        val df =   Seq(
            Array("122abc","223cde","334vbn","445das"),
            Array("221bca","321dsa"),
            Array("231dsa","653asd","698poq","897qwa")
          ).toDF("arr")
          val df2 = df.withColumn("key", 'arr(0)).withColumn("values",explode('arr)).filter('key =!= 'values).drop('arr).withColumn("tuple",struct('key,'values))
          df2.show(false)
          df2.rdd.map( x => Row( (x(0),x(1)) )).collect.foreach(println)
      

      输出:

      +------+------+---------------+
      |key   |values|tuple          |
      +------+------+---------------+
      |122abc|223cde|[122abc,223cde]|
      |122abc|334vbn|[122abc,334vbn]|
      |122abc|445das|[122abc,445das]|
      |221bca|321dsa|[221bca,321dsa]|
      |231dsa|653asd|[231dsa,653asd]|
      |231dsa|698poq|[231dsa,698poq]|
      |231dsa|897qwa|[231dsa,897qwa]|
      +------+------+---------------+
      
      
      [(122abc,223cde)]
      [(122abc,334vbn)]
      [(122abc,445das)]
      [(221bca,321dsa)]
      [(231dsa,653asd)]
      [(231dsa,698poq)]
      [(231dsa,897qwa)]
      

      更新1:

      使用配对的rdd

      val df =   Seq(
        Array("122abc","223cde","334vbn","445das"),
        Array("221bca","321dsa"),
        Array("231dsa","653asd","698poq","897qwa")
      ).toDF("arr")
      import scala.collection.mutable._
      val rdd1 = df.rdd.map( x => { val y = x.getAs[mutable.WrappedArray[String]]("arr")(0); (y,x)} )
      val pair = new PairRDDFunctions(rdd1)
      pair.flatMapValues( x => x.getAs[mutable.WrappedArray[String]]("arr") )
          .filter( x=> x._1 != x._2)
          .collect.foreach(println)
      

      结果:

      (122abc,223cde)
      (122abc,334vbn)
      (122abc,445das)
      (221bca,321dsa)
      (231dsa,653asd)
      (231dsa,698poq)
      (231dsa,897qwa)
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2018-05-23
        • 2016-03-10
        • 1970-01-01
        相关资源
        最近更新 更多