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

如何为Spark DataFrame添加嵌套Struct类型列?

问题:动态生成嵌套Struct类型的Items列

原始DataFrame信息

Schema

root
 |-- Id: string (nullable = true)
 |-- Name: string (nullable = true)

示例数据

+------+------+
| Id   | Name |
+------+------+
|  1   |  'A' |
+------+------+
|  2   |  'B' |
+------+------+

需求目标

需要新增一个Items列,结构为{Id:{"Id":Id,"Name":Name}},最终预期的Schema和数据如下:

预期最终Schema

root
 |-- Id: string (nullable = true)
 |-- Name: string (nullable = true)
 |-- Items: struct (nullable = true)
 |    |-- 1: struct (nullable = true)
 |    |    |-- Id: string (nullable = true)
 |    |    |-- Name: string (nullable = true)
 |    |-- 2: struct (nullable = true)
 |    |    |-- Id: string (nullable = true)
 |    |    |-- Name: string (nullable = true)

预期最终数据

+------+------+-----------------------------+
| Id   | Name | Items                       |
+------+------+-----------------------------+
|  1   |  'A' |  {1:{"Id":1, "Name": "A"}}  |
+------+------+-----------------------------+
|  2   |  'B' |  {2:{"Id":2, "Name": "B"}}  |
+------+------+-----------------------------+

尝试的代码及问题

使用以下代码生成Items列:

transform_expr_items = """struct(Id, struct(Id as Id,
                                      Name as Name
                                ))"""

df_tmp = df_test.withColumn("Items", expr(transform_expr_items))

得到的Schema不符合预期:

root
 |-- Id: string (nullable = true)
 |-- Name: string (nullable = true)
 |-- Items: struct (nullable = true)
 |    |-- Id: string (nullable = true)
 |    |-- col2: struct (nullable = false)
 |    |    |-- Id: string (nullable = true)
 |    |    |-- Name: string (nullable = true)

解决方案

Spark的Struct类型字段名是固定的编译时属性,无法根据每行的动态值生成不同的字段名,因此你想要的结构更适合用Map类型实现,既能达到预期的数据展示效果,又符合Spark的类型系统规范。

推荐实现代码

方式1:使用Spark函数API

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

df_tmp = df_test.withColumn(
    "Items",
    create_map(col("Id"), struct(col("Id").alias("Id"), col("Name").alias("Name")))
)

方式2:使用表达式字符串

from pyspark.sql.functions import expr

transform_expr_items = """create_map(Id, struct(Id as Id, Name as Name))"""
df_tmp = df_test.withColumn("Items", expr(transform_expr_items))

最终效果

生成的Schema如下:

root
 |-- Id: string (nullable = true)
 |-- Name: string (nullable = true)
 |-- Items: map (nullable = false)
 |    |-- key: string
 |    |-- value: struct (nullable = false)
 |    |    |-- Id: string (nullable = true)
 |    |    |-- Name: string (nullable = true)

数据展示与你预期完全一致:

+------+------+-----------------------------+
| Id   | Name | Items                       |
+------+------+-----------------------------+
|  1   |  'A' |  {1:{"Id":1, "Name": "A"}}  |
+------+------+-----------------------------+
|  2   |  'B' |  {2:{"Id":2, "Name": "B"}}  |
+------+------+-----------------------------+

关于Struct实现的说明

如果一定要用Struct类型,只能提前枚举所有可能的Id值,手动构建包含所有字段的Struct,但这种方式灵活性极差,新增Id时需要修改代码,示例如下(仅适用于已知所有Id的场景):

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

df_tmp = df_test.withColumn(
    "Items",
    struct(
        when(col("Id") == "1", struct(col("Id"), col("Name"))).alias("1"),
        when(col("Id") == "2", struct(col("Id"), col("Name"))).alias("2")
    )
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 22:07:45