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
相关产品推荐
相关产品推荐

