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

PySpark展开array<map<string,string>>列并提取数据生成统计宽表

最优实现方案

核心思路使用PySpark原生结构化数据处理能力,避免正则提取的不稳定问题,整体分为JSON解析、数组展开、行转列三个步骤,性能和可维护性更强:

  • 首先导入依赖的函数和类型定义
from pyspark.sql import functions as F
from pyspark.sql.types import ArrayType, StructType, StructField, StringType, IntegerType
  • 定义数组结构的Schema,将字符串类型的统计值转为结构化数组
    由于你原始employment_status的map值为字符串类型,需要先通过from_json将存储的统计数组解析为结构化对象数组,若你当前sdf4的value列已经是数组类型可跳过解析步骤直接展开:
# 定义就业统计数组的结构
status_schema = ArrayType(
    StructType([
        StructField("x", StringType(), nullable=True),
        StructField("y", IntegerType(), nullable=True)
    ])
)

# 解析并展开数组,得到每行一个就业类型的长表
sdf_long = sdf4 \
    .withColumn("status_arr", F.from_json("value", status_schema)) \
    .select("zipcode", F.explode("status_arr").alias("item")) \
    .select(
        "zipcode",
        F.col("item.x").alias("emp_type"),
        F.col("item.y").alias("count")
    )
  • 行转列得到宽表
    使用pivot函数将就业类型转为列,直接得到按邮编分组的结构化宽表,支持自定义列名:
sdf_wide = sdf_long \
    .groupBy("zipcode") \
    .pivot("emp_type") \
    .agg(F.first("count")) \
    # 自定义列名,可按需求调整
    .withColumnRenamed("Full-time", "full_time_emp_count") \
    .withColumnRenamed("Part-time", "part_time_emp_count") \
    .withColumnRenamed("No Earnings", "no_earnings_count")

最终输出样例:

+-------+-------------------+-------------------+------------------+
|zipcode|full_time_emp_count|part_time_emp_count|no_earnings_count |
+-------+-------------------+-------------------+------------------+
|95678  |13348              |8918               |9972              |
|95679  |0                  |29                 |0                 |
+-------+-------------------+-------------------+------------------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 12:06:03