如何将含字典数组的列转为多列?处理未知内容及嵌套结构
问题描述
我有一个DataFrame,其中data和modules两列包含字典数组,尝试将其展开为多列但未成功。
DataFrame Schema
|-- data: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- name: string (nullable = true) | | |-- value: string (nullable = true) | | |-- origin: string (nullable = true) |-- createdAt: long (nullable = true) |-- modules: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- name: string (nullable = true) | | |-- detected: struct (nullable = true) | | | |-- name: string (nullable = true) | | | |-- id: integer (nullable = true) | | | |-- reason: string (nullable = true) | | | |-- state: string (nullable = true) | | | |-- score: integer (nullable = true) | | | |-- level: string (nullable = true) |-- id: string (nullable = true)
数据样例
+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ |data |createdAt |modules |id | +----------------------------------------------------------------------+-------------+-------------------------------------------------------------------------------------------------------------------------------------------+------------------------------------+ |[{data_point_1, false, METADATA}, {data_point_2, some_string, DEVICE}]|1678148428468|[{ANOTHER_DUMMY, {null, null, null, null, null, null}}, {DUMMY, {dummy_user_agent, 1, Rule for integration tests, OPERATIONAL, 500, HIGH}}]|70ef58bf-b160-4abd-97c1-aa4780e74e1b| |[{data_point_1, false, METADATA}, {data_point_3, 0, USER}] |1678148428495|[{ANOTHER_DUMMY, {null, null, null, null, null, null}}, {DUMMY, {null, null, null, null, null, null}}] |6ab33e95-dd94-4c95-b00f-edfe97d6f3d1| +---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
期望输出
需要将其转换为如下宽表格式:
+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ |data_point_1|data_point_1_origin|data_point_2|data_point_2_origin|data_point_3|data_point_3_origin|createdAt |ANOTHER_DUMMY_name|ANOTHER_DUMMY_id|ANOTHER_DUMMY_reason|ANOTHER_DUMMY_state|ANOTHER_DUMMY_score|ANOTHER_DUMMY_level|DUMMY_name |DUMMY_id|DUMMY_reason |DUMMY_state|DUMMY_score|DUMMY_level|id | +------------+-------------------+------------+-------------------+------------+-------------------+-------------+------------------+----------------+--------------------+-------------------+-------------------+-------------------+----------------+--------+--------------------------+-----------+-----------+-----------+------------------------------------+ |false |METADATA |some_string |DEVICE |null |null |1678148428468|null |null |null |null |null |null |dummy_user_agent|1 |Rule for integration tests|OPERATIONAL|500 |HIGH |70ef58bf-b160-4abd-97c1-aa4780e74e1b| |false |METADATA |null |null |0 |USER |1678148428495|null |null |null |null |null |null |null |null |null |null |null |null |6ab33e95-dd94-4c95-b00f-edfe97d6f3d1| +-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
要求:未知data列字典内容,需处理modules列嵌套结构。
解决方案
以下以PySpark为例,分步骤处理data和modules列的展开:
1. 处理data列:动态展开键值对
data列是包含name/value/origin的结构体数组,且每行name不固定,需先转成键值对形式再动态生成列:
from pyspark.sql import functions as F # 拆分data数组为多行,每行对应一个结构体 data_exploded = df.withColumn("data_item", F.explode("data")) # 按id、createdAt分组,以data_item.name为列名聚合value和origin data_pivoted = data_exploded.groupBy("id", "createdAt", "modules") \ .pivot("data_item.name") \ .agg( F.first("data_item.value").alias("value"), F.first("data_item.origin").alias("origin") ) # 重命名列,将xx_value改为xx,xx_origin保留原格式 final_data_df = data_pivoted.select( "id", "createdAt", "modules", *[F.col(f"{col}_value").alias(col) for col in [c.split("_value")[0] for c in data_pivoted.columns if "_value" in c]], *[F.col(f"{col}_origin").alias(f"{col}_origin") for col in [c.split("_origin")[0] for c in data_pivoted.columns if "_origin" in c]] )
2. 处理modules列:展开嵌套结构体并动态命名列
modules列是包含name和嵌套detected结构体的数组,需先展开数组,再提取detected字段并按模块名_字段名格式命名:
# 展开modules数组 modules_exploded = final_data_df.withColumn("module_item", F.explode("modules")) # 提取module名称和detected的所有字段,组合成新列名 detected_fields = ["name", "id", "reason", "state", "score", "level"] modules_flattened = modules_exploded.select( "id", "createdAt", *[col for col in final_data_df.columns if col not in ["modules"]], *[F.col(f"module_item.detected.{field}").alias(f"{F.col('module_item.name')}_{field}") for field in detected_fields] ) # 动态获取所有唯一的module名称(适配module不固定的场景) module_names = [row["module_name"] for row in modules_exploded.select(F.col("module_item.name").alias("module_name")).distinct().collect()] # 按id和createdAt分组,聚合取第一个非空值(合并同一id的多个module行) final_df = modules_flattened.groupBy("id", "createdAt", *[col for col in final_data_df.columns if col not in ["modules", "id", "createdAt"]]) \ .agg( *[F.first(f"{module}_{field}", ignorenulls=True).alias(f"{module}_{field}") for module in module_names for field in detected_fields] )
3. 调整列顺序(可选)
如果需要和期望输出的列顺序一致,可手动指定列顺序:
desired_columns = [ "data_point_1", "data_point_1_origin", "data_point_2", "data_point_2_origin", "data_point_3", "data_point_3_origin", "createdAt", "ANOTHER_DUMMY_name", "ANOTHER_DUMMY_id", "ANOTHER_DUMMY_reason", "ANOTHER_DUMMY_state", "ANOTHER_DUMMY_score", "ANOTHER_DUMMY_level", "DUMMY_name", "DUMMY_id", "DUMMY_reason", "DUMMY_state", "DUMMY_score", "DUMMY_level", "id" ] final_df = final_df.select(desired_columns)
内容的提问来源于stack exchange,提问作者Ema Il
相关产品推荐
相关产品推荐

