【问题标题】:Pivot spark dataframe array of kv pairs into individual columns将 kv 对的 spark 数据帧数组转换为单独的列
【发布时间】:2019-08-12 18:23:48
【问题描述】:

我有以下架构:

root
 |-- id: string (nullable = true)
 |-- date: timestamp (nullable = true)
 |-- config: struct (nullable = true)
 |    |-- entry: array (nullable = true)
 |    |    |-- element: struct (containsNull = true)
 |    |    |    |-- key: string (nullable = true)
 |    |    |    |-- value: string (nullable = true)

数组中的键值对 (k1,k2,k3) 不会超过 3 个,我想将每个键的值放入自己的列中,而相应的数据将来自同一列的值kv 对。

+--------+----------+----------+----------+---------+
|id      |date      |k1        |k2        |k3       |
+--------+----------+----------+----------+---------+
|    id1 |2019-08-12|id1-v1    |id1-v2    |id1-v3   |
|    id2 |2019-08-12|id2-v1    |id2-v2    |id2-v3   |
+--------+----------+----------+----------+---------+

到目前为止,我尝试过这样的事情:

sourceDF.filter($"someColumn".contains("SOME_STRING"))
      .select($"id", $"date", $"config.entry" as "kvpairs")
      .withColumn($"kvpairs".getItem(0).getField("key").toString(), $"kvpairs".getItem(0).getField("value"))
      .withColumn($"kvpairs".getItem(1).getField("key").toString(), $"kvpairs".getItem(1).getField("value"))
      .withColumn($"kvpairs".getItem(2).getField("key").toString(), $"kvpairs".getItem(2).getField("value"))

但在这种情况下,列名显示为kvpairs[0][key]kvpairs[1][key]kvpairs[2][key],如下所示:

+--------+----------+---------------+---------------+---------------+
|id      |date      |kvpairs[0][key]|kvpairs[1][key]|kvpairs[2][key]|
+--------+----------+---------------+---------------+---------------+
|    id1 |2019-08-12|    id1-v1     |    id1-v2     |   id1-v3      |
|    id2 |2019-08-12|    id2-v1     |    id2-v2     |   id2-v3      |
+--------+----------+---------------+---------------+---------------+

两个问题:

  • 我的方法对吗?有没有更好更简单的方法来解决这个问题 这样我每个数组就得到一行,3 kv 对作为 3 列?我想处理 kv 对的顺序可能不同的情况。
  • 如果上述方法没问题,如何将列名别名为数组中“key”元素的数据?

【问题讨论】:

  • 您的方法可能是最佳的。使用alias 重命名列。另外-您标记了 pyspark,但您的代码看起来像 scala
  • 如何以编程方式在 withColumn 中使用别名,使得列名是与键对应的数据。例如: [[a,1],[b,2],[c,3]] 。我希望列名是 a,b,c,值分别是 1,2,3。一个示例 sn-p 可能会有所帮助。

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


【解决方案1】:

同时使用多个withColumngetItem 将不起作用,因为kv 对的顺序可能不同。你可以做的是分解数组,然后使用pivot,如下所示:

sourceDF.filter($"someColumn".contains("SOME_STRING"))
  .select($"id", $"date", explode($"config.entry") as "exploded")
  .select($"id", $"date", $"exploded.*")
  .groupBy("id", "date")
  .pivot("key")
  .agg(first("value"))

这里在聚合中使用first 假设每个键都有一个值。否则可以使用collect_listcollect_set

结果:

+---+----------+------+------+------+
|id |date      |k1    |k2    |k2    |
+---+----------+------+------+------+
|id1|2019-08-12|id1-v1|id1-v2|id1-v3|
|id2|2019-08-12|id2-v1|id2-v2|id2-v3|
+---+----------+------+------+------+

【讨论】:

  • 谢谢!这就像一个魅力!如果排序保持不变,我假设 withColumn 方法可以正常工作?
  • @arjunj:是的,这是正确的。如果是这种情况,那么首先从数据框中的第一行中找到键名,然后在 withColumn 中使用 alias/as 应用这些键名将起作用。这样标题就正确了。
猜你喜欢
  • 2018-07-16
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-09-27
  • 2019-01-16
相关资源
最近更新 更多