【问题标题】:Reorganize Pyspark dataframe: Create new column using row element重组 Pyspark 数据框:使用行元素创建新列
【发布时间】:2020-10-23 11:59:41
【问题描述】:

我正在尝试将具有此结构的文档映射到数据框。

   root
     |-- Id: "a1"
     |-- Type: "Work"
     |-- Tag: Array
     |    |--0: Object 
     |    |   |-- Tag.name : "passHolder"
     |    |   |-- Tag.value : "Jack Ryan"
     |    |   |-- Tag.stat : "verified"
     |    |-- 1: Object
     |    |   |-- Tag.name : "passNum"
     |    |   |-- Tag.value : "1234"
     |    |   |-- Tag.stat : "unverified"
     |-- version: 1.5
                

通过使用explode_outer 分解数组、展平结构并使用.col + alias 重命名,数据框将如下所示:

df = df.withColumn("Tag",F.explode_outer("Tag"))
df = df.select(col("*"), 
       .col("Tag.name").alias("Tag_name"),
       .col("Tag.value").alias("Tag_value"),
       .col("Tag.stat").alias("Tag_stat")).drop("Tag")

+--+----+-----------+-----------+---------+---------+
|Id|Type| Tag_name  | Tag_value |Tag_stat | version |
+--+----+-----------+-----------+---------+---------+
 a1 Work  passHolder  Jack Ryan  verified       1.5
 a1 Work  passNum      1234      unverified     1.5

我正在尝试重新组织 df 结构,使其更易于查询,方法是将某些行元素作为列名并用相关值填充它。 任何人都可以帮助提供达到所需输出格式所需的指针/步骤,如下所示?非常感谢您的建议。

目标格式:

+--+----+-----------------+-----------------+-------------+------------+--------+
|Id|Type| Tag_passHolder  | passHolder_stat | Tag_passNum |passNum_stat||version|
+--+----+-----------------+-----------------+-------------+------------+--------+
 a1 Work   Jack Ryan          verified           1234       unverified     1.5   

【问题讨论】:

  • 听起来像一个简单的连接。你试过了吗?
  • @Steven 我没有。你能帮助扩展如何使用连接来完成目标 sparkdf 吗?

标签: python pyspark apache-spark-sql


【解决方案1】:

根据您显示的输出 df,我会做这样的事情:

from pyspark.sql import functions as F

passholder_df = df.select(
    "ID",
    "Type",
    F.col("Tag_value").alias("Tag_passHolder"),
    F.col("Tag_stat").alias("passHolder_stat"),
    "version",
).where("Tag_name = 'passHolder'")

passnum_df = df.select(
    "ID",
    "Type",
    F.col("Tag_value").alias("Tag_passNum"),
    F.col("Tag_stat").alias("passNum_stat"),
    "version",
).where("Tag_name = 'passNum'")

passholder_df.join(passnum_df, on=["ID", "Type", "version"], how="full")

根据您的业务规则,您可能需要对连接条件进行一些处理。

【讨论】:

  • 我在这里得到了你的方法@Steven,谢谢。在Tag 数组是动态的并且可以拥有更多具有相同结构的对象的情况下,是否可以(递归地?)将Tag.name 值检索为新列并为每个对象使用Tag.value 填充它?
猜你喜欢
  • 1970-01-01
  • 2021-10-07
  • 1970-01-01
  • 2019-05-11
  • 2022-06-21
  • 2018-09-29
  • 1970-01-01
  • 1970-01-01
  • 2020-11-07
相关资源
最近更新 更多