【问题标题】:Zip and Explode multiple Columns in Spark SQL Dataframe在 Spark SQL Dataframe 中压缩和分解多个列
【发布时间】:2023-03-25 07:00:01
【问题描述】:

我有一个如下结构的数据框:

A: Array[String]   | B: Array[String] | [ ... multiple other columns ...]
=========================================================================
[A, B, C, D]       | [1, 2, 3, 4]     | [ ... array with 4 elements ... ]
[E, F, G, H, I]    | [5, 6, 7, 8, 9]  | [ ... array with 5 elements ... ]
[J]                | [10]             | [ ... array with 1 element ...  ]

我想写一个UDF,那个

  1. 压缩 DF 中每列第 i 个位置的元素
  2. 在每个压缩元组上分解 DF

生成的列应如下所示:

ZippedAndExploded: Array[String]
=================================
[A, 1, ...]
[B, 2, ...]
[C, 3, ...]
[D, 4, ...]
[E, 5, ...]
[F, 6, ...]
[G, 7, ...]
[H, 8, ...]
[I, 9, ...]
[J, 10, ...]

目前我正在对这样的 UDF 使用多重调用(每个列名一个,列名列表在运行时之前收集):

val myudf6 = udf((xa:Seq[Seq[String]],xb:Seq[String]) => {
  xa.indices.map(i => {
    xa(i) :+ xb(i) // Add one element to the zip column
  })
})

val allColumnNames = df.columns.filter(...)    

for (columnName <- allColumnNames) {
  df = df.withColumn("zipped", myudf8(df("zipped"), df(columnName))
}
df = df.explode("zipped")

由于数据框可以有数百列,withColumn 的这种迭代调用似乎需要很长时间。

问题:这可能与一个 UDF 和一个 DF.withColumn(...) 调用有关吗?

重要提示:UDF 应压缩动态数量的列(在运行时读取)。

【问题讨论】:

  • 知道如何在 PySpark 中进行操作吗?

标签: apache-spark apache-spark-sql user-defined-functions apache-spark-dataset


【解决方案1】:

使用将可变数量的列作为输入的UDF。这可以通过数组数组来完成(假设类型相同)。由于您有一个数组数组,因此可以使用transpose,这将获得与将列表压缩在一起的相同结果。然后可以分解生成的数组。

val array_zip_udf = udf((cols: Seq[Seq[String]]) => {
  cols.transpose
})

val allColumnNames = df.columns.filter(...).map(col)
val df2 = df.withColumn("exploded", explode(array_zip_udf(array(allColumnNames: _*))))

请注意,在 Spark 2.4+ 中,可以使用 arrays_zip 代替 UDF

val df2 = df.withColumn("exploded", explode(arrays_zip(allColumnNames: _*)))

【讨论】:

  • 使用@Shaido 的代码解决了它(将最后一行更改为val df2 = df.withColumn("exploded", explode(array_zip_udf(array(columnNames.head, columnNames.tail : _*))))
  • 非常感谢,您为我节省了很多时间!
  • @D.Müller:乐于助人 :) 我固定回答一点以避免head/tail 部分,这可以通过在列名中添加map(col) 来完成。 (我忘了array 只接受Seq[Column] 而不是Seq[String]。)
【解决方案2】:

如果您知道并确定数组中值的数量,以下可能是更简单的解决方案之一

select A[0], B[0]..... from your_table
union all
select A[1], B[1]..... from your_table
union all
select A[2], B[2]..... from your_table
union all
select A[3], B[3]..... from your_table

【讨论】:

  • 感谢您的快速回复。我调整了我的问题,因为数据框可以包含不同大小的数组(按行)。我想,这个解决方案不适用于此类数据,对吧?很抱歉问题的迟来的变化。
猜你喜欢
  • 2015-12-29
  • 2019-06-14
  • 2016-01-18
  • 2021-05-25
  • 2018-12-15
  • 2010-10-15
  • 1970-01-01
  • 1970-01-01
  • 2018-01-09
相关资源
最近更新 更多