You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在PySpark中将DataFrame的数组类型列转换为结构体

PySpark数组类型列转结构体类型列

需求说明

将DataFrame中的documents和contacts数组列转换为结构体列,结构体包含element字段,对应原数组中的单个元素。

输入Schema

root
 |-- id: string (nullable = true)
 |-- description: string (nullable = true)
 |-- documents: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- id: string (nullable = true)
 |    |    |-- doc_name: string (nullable = true)
 |    |    |-- obligations: struct (containsNull = true)
 |-- contacts: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- id: string (nullable = true)
 |    |    |-- contact_first_name: string (nullable = true)
 |    |    |-- contact_last_name: string (nullable = true)

输入数据

{
   "id":"123",
   "description": "agreement",
   "documents":[
     {
       "id":"doc_id_1",
       "doc_name":"doc_name_1",
       "obligations":{}
     }
   ],
   "contacts":[
    {
      "id":"contact_id_1",
      "contact_first_name":"John",
      "contact_last_name":"Doe"
    }
  ]
}

解决方案代码

1. 创建示例DataFrame

from pyspark.sql import SparkSession
from pyspark.sql.functions import struct, col

spark = SparkSession.builder.appName("ArrayToStructConversion").getOrCreate()

input_data = [
    {
        "id": "123",
        "description": "agreement",
        "documents": [{"id": "doc_id_1", "doc_name": "doc_name_1", "obligations": {}}],
        "contacts": [{"id": "contact_id_1", "contact_first_name": "John", "contact_last_name": "Doe"}]
    }
]

df = spark.createDataFrame(input_data)

2. 执行数组转结构体转换

针对示例中数组仅含单个元素的场景,直接取数组首个元素包装为结构体:

transformed_df = df.withColumn(
    "documents",
    struct(col("documents").getItem(0).alias("element"))
).withColumn(
    "contacts",
    struct(col("contacts").getItem(0).alias("element"))
)

3. 处理空数组/Null场景(可选)

如果数组可能为空或为Null,添加空值兼容逻辑:

from pyspark.sql.functions import when

transformed_df = df.withColumn(
    "documents",
    when(
        col("documents").isNull() | (col("documents").size() == 0),
        struct(col("documents").getItem(0).alias("element"))
    ).otherwise(
        struct(col("documents").getItem(0).alias("element"))
    )
).withColumn(
    "contacts",
    when(
        col("contacts").isNull() | (col("contacts").size() == 0),
        struct(col("contacts").getItem(0).alias("element"))
    ).otherwise(
        struct(col("contacts").getItem(0).alias("element"))
    )
)

验证结果

查看转换后的Schema:

transformed_df.printSchema()

输出与期望目标Schema一致:

root
 |-- id: string (nullable = true)
 |-- description: string (nullable = true)
 |-- documents: struct (containsNull = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- id: string (nullable = true)
 |    |    |-- doc_name: string (nullable = true)
 |    |    |-- obligations: struct (containsNull = true)
 |-- contacts: struct (containsNull = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- id: string (nullable = true)
 |    |    |-- contact_first_name: string (nullable = true)
 |    |    |-- contact_last_name: string (nullable = true)

查看转换后的数据:

transformed_df.toJSON().collect()

输出符合期望格式:

[
  {
    "id": "123",
    "description": "agreement",
    "documents": {
      "element": {
        "id": "doc_id_1",
        "doc_name": "doc_name_1",
        "obligations": {}
      }
    },
    "contacts": {
      "element": {
        "id": "contact_id_1",
        "contact_first_name": "John",
        "contact_last_name": "Doe"
      }
    }
  }
]

内容的提问来源于stack exchange,提问作者bda

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 16:55:31