【问题标题】:Spark Dataframe implementation similar to Oracle's LISTAGG function - Unable to Order with in the groupSpark Dataframe 实现类似于 Oracle 的 LISTAGG 函数 - 无法在组中订购
【发布时间】:2018-12-15 15:31:47
【问题描述】:

我想实现一个类似于Oracle的LISTAGG函数的函数。

等价的oracle代码是

select KEY,
listagg(CODE, '-') within group (order by DATE) as CODE
from demo_table
group by KEY

这是我的 spark scala 数据框实现,但无法对每个组中的值进行排序。

输入:

val values = List(List("66", "PL", "2016-11-01"), List("66", "PL", "2016-12-01"),
  List("67", "JL", "2016-12-01"), List("67", "JL", "2016-11-01"), List("67", "PL", "2016-10-01"), List("67", "PO", "2016-09-01"), List("67", "JL", "2016-08-01"),
  List("68", "PL", "2016-12-01"), List("68", "JO", "2016-11-01"))
  .map(row => (row(0), row(1), row(2)))

val df = values.toDF("KEY", "CODE", "DATE")

df.show()

+---+----+----------+
|KEY|CODE|      DATE|
+---+----+----------+
| 66|  PL|2016-11-01|
| 66|  PL|2016-12-01|----- group 1
| 67|  JL|2016-12-01|
| 67|  JL|2016-11-01|
| 67|  PL|2016-10-01|
| 67|  PO|2016-09-01|
| 67|  JL|2016-08-01|----- group 2
| 68|  PL|2016-12-01|
| 68|  JO|2016-11-01|----- group 3
+---+----+----------+

udf 实现:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.functions.udf

val listAgg = udf((xs: Seq[String]) => xs.mkString("-"))

df.groupBy("KEY")
  .agg(listAgg(collect_list("CODE")).alias("CODE"))
  .show(false)

+---+--------------+
|KEY|CODE          |
+---+--------------+
|68 |PL-JO         |
|67 |JL-JL-PL-PO-JL|
|66 |PL-PL         |
+---+--------------+

预期输出:- 按日期排序

+---+--------------+
|KEY|CODE          |
+---+--------------+
|68 |JO-PL         |
|67 |JL-PO-PL-JL-JL|
|66 |PL-PL         |
+---+--------------+

【问题讨论】:

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


    【解决方案1】:

    使用struct 内置函数 组合CODEDATE 列,并在collect_list 聚合函数中使用该新结构列。在udf 函数中按日期排序收集代码 作为- 分隔字符串

    import org.apache.spark.sql.functions._
    def sortAndStringUdf = udf((codeDate: Seq[Row])=> codeDate.sortBy(row => row.getAs[Long]("DATE")).map(row => row.getAs[String]("CODE")).mkString("-"))
    
    df.withColumn("codeDate", struct(col("CODE"), col("DATE").cast("timestamp").cast("long").as("DATE")))
          .groupBy("KEY").agg(sortAndStringUdf(collect_list("codeDate")).as("CODE"))
    

    这应该给你

    +---+--------------+
    |KEY|          CODE|
    +---+--------------+
    | 68|         JO-PL|
    | 67|JL-PO-PL-JL-JL|
    | 66|         PL-PL|
    +---+--------------+
    

    希望回答对你有帮助

    更新

    我相信这会比使用udf 函数更快

    df.withColumn("codeDate", struct(col("DATE").cast("timestamp").cast("long").as("DATE"), col("CODE")))
      .groupBy("KEY")
      .agg(concat_ws("-", expr("sort_array(collect_list(codeDate)).CODE")).alias("CODE"))
      .show(false)
    

    这应该会给你与上面相同的结果

    【讨论】:

      猜你喜欢
      • 2017-07-05
      • 2022-12-14
      • 2017-10-18
      • 2016-07-25
      • 1970-01-01
      • 2021-11-30
      • 1970-01-01
      • 2016-09-02
      相关资源
      最近更新 更多