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

Spark中如何从嵌套Struct列提取key4的键值对

从Spark嵌套Struct列提取指定子结构的键值对

错误原因

你之前调用items()报错是因为Spark的Struct类型在Python UDF中对应pyspark.sql.types.Row对象,它不是普通Python字典,没有items()方法,直接调用会触发AttributeError。

解决方案

方法1:使用Spark内置函数(推荐,性能更优)

Spark提供了内置函数可以直接处理Struct类型,无需编写UDF,适合大数据场景。

假设你的数据结构类似:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

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

# 模拟数据与schema
data = [
    (("val1", "val2", "val3", {"subkey1": 1, "subkey2": 2}),)
]
schema = StructType([
    StructField("column_1", StructType([
        StructField("key1", StringType()),
        StructField("key2", StringType()),
        StructField("key3", StringType()),
        StructField("key4", StructType([
            StructField("subkey1", IntegerType()),
            StructField("subkey2", IntegerType())
        ]))
    ]))
])
df = spark.createDataFrame(data, schema)

提取key4的键值对步骤:

  1. 先获取key4子结构的所有字段名
  2. 用struct生成键值对结构体,再用array组合成数组,最后用map_from_entries转成Map类型(或直接保留数组)

代码实现:

from pyspark.sql.functions import lit, struct, array, map_from_entries, col

# 获取key4的所有字段名
key4_fields = [field.name for field in df.schema["column_1"].dataType["key4"].fields]

# 生成包含所有键值对的Map列
df = df.withColumn(
    "key4_key_value",
    map_from_entries(
        array(*[struct(lit(field), col(f"column_1.key4.{field}")) for field in key4_fields])
    )
)

# 查看结果
df.select("key4_key_value").show(truncate=False)

输出结果:

+---------------------------+
|key4_key_value             |
+---------------------------+
|{subkey1 -> 1, subkey2 -> 2}|
+---------------------------+

如果需要数组形式的键值对(每个元素是(key, value)结构体),直接使用array部分即可:

df = df.withColumn(
    "key4_key_value_array",
    array(*[struct(lit(field), col(f"column_1.key4.{field}")) for field in key4_fields])
)

方法2:修正UDF写法

如果一定要用UDF,需要先将Row对象转为Python字典,再调用items()方法。

代码实现:

from pyspark.sql.functions import udf
from pyspark.sql.types import ArrayType, StructType, StructField, StringType, IntegerType

# 定义UDF的返回类型:数组,每个元素是包含key和value的结构体
kv_schema = ArrayType(StructType([
    StructField("key", StringType()),
    StructField("value", IntegerType())  # 若key4字段类型多样,可改为StringType或AnyType
]))

@udf(kv_schema)
def extract_key4_kv(row):
    # 将Row对象转为字典
    key4_dict = row.key4.asDict()
    # 返回键值对列表
    return list(key4_dict.items())

# 应用UDF
df = df.withColumn("key4_key_value", extract_key4_kv(col("column_1")))

# 查看结果
df.select("key4_key_value").show(truncate=False)

输出结果:

+-------------------------+
|key4_key_value           |
+-------------------------+
|[{subkey1, 1}, {subkey2, 2}]|
+-------------------------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 18:15:38