在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
相关产品推荐
相关产品推荐

