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的键值对步骤:
- 先获取
key4子结构的所有字段名 - 用
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
相关产品推荐
相关产品推荐

