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

PySpark如何将key:value字符串数组转为struct并提取指定字段值

PySpark实现键值数组字段结构化提取

实现思路

  • 先将data_zone_array中每个key:value格式的字符串拆分为独立的键值对结构
  • 按规则提取预设字段name(唯一字符串)和surname(字符串数组)
  • 剩余非预设字段先聚合为「键对应值数组」的映射结构,再转为要求的struct格式

完整代码示例

from pyspark.sql import functions as F
from pyspark.sql.types import MapType, ArrayType, StringType

# 假设原始DataFrame为df
# 1. 拆分每个元素为键值对结构,限制按第一个冒号拆分避免值中含冒号出错
df = df.withColumn("kv_pairs", F.transform(
    "data_zone_array",
    lambda x: F.struct(
        F.split(x, ":", 2)[0].alias("key"),
        F.split(x, ":", 2)[1].alias("value")
    )
))

# 2. 提取预设字段name:唯一值直接取第一个匹配结果
df = df.withColumn("name", F.element_at(
    F.filter("kv_pairs", lambda x: x.key == "name").value,
    1
))

# 3. 提取预设字段surname:直接返回所有匹配值的数组
df = df.withColumn("surname", F.filter(
    "kv_pairs", lambda x: x.key == "surname"
).value)

# 4. 过滤出非预设键的键值对
df = df.withColumn("other_kv", F.filter(
    "kv_pairs", lambda x: ~x.key.isin("name", "surname")
))

# 5. 行内聚合非预设键为 map<key, array<value>> 结构
df = df.withColumn("other_map", F.aggregate(
    "other_kv",
    F.create_map().cast(MapType(StringType(), ArrayType(StringType()))),
    lambda acc, curr: F.map_concat(
        acc,
        F.create_map(
            curr.key,
            F.when(
                F.map_contains_key(acc, curr.key),
                F.concat(acc[curr.key], F.array(curr.value))
            ).otherwise(F.array(curr.value))
        )
    )
))

# 6. 批处理场景:收集所有出现过的非预设键,将map转为struct
all_other_keys = [
    row["key"] 
    for row in df.select(F.explode("other_kv.key")).distinct().collect()
]
df = df.withColumn("other_attributes", F.struct(*[
    F.col("other_map")[k].alias(k) for k in all_other_keys
]))

# 7. 删除中间字段,得到最终结果
df_final = df.drop("data_zone_array", "kv_pairs", "other_kv", "other_map")

# 验证结果
df_final.printSchema()
df_final.show(truncate=False)

注意事项

  • 流处理场景无法提前枚举所有非预设键时,可直接保留other_map(map类型)字段,无需转为struct,使用时可直接通过key取值
  • 若需要surname或其他属性值按字母序排列,可在提取时增加array_sort函数处理
  • 代码适配Spark 3.0及以上版本,低版本Spark可通过自定义UDF实现相同逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 03:45:05