【问题标题】:Select the last element of an Array in a DataFrame在 DataFrame 中选择 Array 的最后一个元素
【发布时间】:2018-11-08 14:50:57
【问题描述】:

我正在处理一个项目,我正在处理一些具有复杂架构/数据结构的嵌套 JSON 日期。基本上我想要做的是过滤掉数据框中的一列,以便我选择数组中的最后一个元素。我完全坚持如何做到这一点。我希望这是有道理的。

以下是我正在尝试完成的示例:

val singersDF = Seq(
  ("beatles", "help,hey,jude"),
  ("romeo", "eres,mia"),
  ("elvis", "this,is,an,example")
).toDF("name", "hit_songs")

val actualDF = singersDF.withColumn(
  "hit_songs",
  split(col("hit_songs"), "\\,")
)

actualDF.show(false)
actualDF.printSchema() 

+-------+-----------------------+
|name   |hit_songs              |
+-------+-----------------------+
|beatles|[help, hey, jude]      |
|romeo  |[eres, mia]            |
|elvis  |[this, is, an, example]|
+-------+-----------------------+
root
 |-- name: string (nullable = true)
 |-- hit_songs: array (nullable = true)
 |    |-- element: string (containsNull = true)

输出的最终目标如下,选择 hit_songs 数组中的最后一个“字符串”。

我不担心架构之后会是什么样子。

+-------+---------+
|name   |hit_songs|
+-------+---------+
|beatles|jude     |
|romeo  |mia      |
|elvis  |example  |
+-------+---------+

【问题讨论】:

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


    【解决方案1】:

    您可以使用size 函数计算所需项在数组中的索引,然后将其作为Column.apply 的参数(显式或隐式)传递:

    import org.apache.spark.sql.functions._
    import spark.implicits._
    
    actualDF.withColumn("hit_songs", $"hit_songs".apply(size($"hit_songs").minus(1)))
    

    或者:

    actualDF.withColumn("hit_songs", $"hit_songs"(size($"hit_songs").minus(1)))
    

    【讨论】:

      【解决方案2】:

      这是一种方法:

      val actualDF = Seq(
        ("beatles", Seq("help", "hey", "jude")),
        ("romeo", Seq("eres", "mia")),
        ("elvis", Seq("this", "is", "an", "example"))
      ).toDF("name", "hit_songs")
      
      import org.apache.spark.sql.functions._
      
      actualDF.withColumn("total_songs", size($"hit_songs")).
        select($"name", $"hit_songs"($"total_songs" - 1).as("last_song"))
      // +-------+---------+
      // |   name|last_song|
      // +-------+---------+
      // |beatles|     jude|
      // |  romeo|      mia|
      // |  elvis|  example|
      // +-------+---------+
      

      【讨论】:

        【解决方案3】:

        您正在寻找 SparkSQL 函数 slice。或者这个PySpark Source

        您在 Scala 中的实现 slice($"hit_songs", -1, 1)(0) 其中-1 是起始位置(最后一个索引),1 是长度,(0) 从恰好包含 1 个元素的结果数组中提取第一个字符串。

        完整示例:

        import org.apache.spark.sql.functions._
        
        val singersDF = Seq(
          ("beatles", "help,hey,jude"),
          ("romeo", "eres,mia"),
          ("elvis", "this,is,an,example")
        ).toDF("name", "hit_songs")
        
        val actualDF = singersDF.withColumn(
          "hit_songs",
          split(col("hit_songs"), "\\,")
        )
        
        val newDF = actualDF.withColumn("last_song", slice($"hit_songs", -1, 1)(0))
        
        display(newDF)
        

        输出:

        +---------+------------------------------+-------------+
        |  name   |          hit_songs           |  last_song  |
        +---------+------------------------------+-------------+
        | beatles | ["help","hey","jude"]        | jude        |
        | romeo   | ["eres","mia"]               | mia         |
        | elvis   | ["this","is","an","example"] | example     |
        +---------+------------------------------+-------------+
        

        【讨论】:

          【解决方案4】:

          spark 2.4+ 开始,您可以使用支持负索引的element_at。正如您在此文档引用中看到的:

          element_at(array, index) - 返回给定(从 1 开始)索引处的数组元素。如果 index

          这样,获取最后一个元素的方法如下:

          import org.apache.spark.sql.functions.element_at
          actualDF.withColumn("hit_songs", element_at($"hit_songs", -1))
          

          可重现的例子:

          首先让我们准备一个带有数组列的示例数据框:

          val columns = Seq("col1")
          val data = Seq((Array(1,2,3)))
          val rdd = spark.sparkContext.parallelize(data)
          val df = rdd.toDF(columns:_*)
          

          看起来像这样:

          scala> df.show()
          +---------+
          |     col1|
          +---------+
          |[1, 2, 3]|
          +---------+
          

          然后,应用element_at 得到最后一个元素如下:

          scala> df.withColumn("last_value", element_at($"col1", -1)).show()
          +---------+----------+
          |     col1|last_value|
          +---------+----------+
          |[1, 2, 3]|         3|
          +---------+----------+
          

          【讨论】:

            【解决方案5】:

            您还可以使用如下 UDF:

            val lastElementUDF = udf((array: Seq[String]) => array.lastOption)
            
            actualDF.withColumn("hit_songs", lastElementUDF($"hit_songs"))
            

            array.lastOption 将返回 NoneSome,如果数组为空,array.last 将抛出异常。

            【讨论】:

              猜你喜欢
              • 2017-11-28
              • 2013-08-18
              • 1970-01-01
              • 2022-06-23
              • 1970-01-01
              • 2012-05-14
              • 1970-01-01
              • 1970-01-01
              • 2011-11-07
              相关资源
              最近更新 更多