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

PySpark如何将DataFrame多列数组按序拆分为独立行

PySpark 多数组列按位置对齐拆分行实现

核心思路

不要对单个数组列单独做explode再做关联,会产生笛卡尔积、还可能出现位置错位问题。直接用数组打包+外爆的方式实现:

  • 用arrays_zip将所有需要拆分的数组列按位置对齐,打包成结构体类型的数组,空数组会被正确识别
  • 用explode_outer替代普通explode,空数组/Null数组的行不会被丢弃,会保留一行对应空值
  • 从外爆得到的结构体中取出各字段,还原成原始列名即可

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

# 初始化Spark会话,生产环境可移除master本地配置
spark = SparkSession.builder.master("local[*]").appName("array_split_row").getOrCreate()

# 构造与示例一致的测试源数据
source_data = [
    ("449", "80212", [], [], [], [], [], []),
    ("449", "80214", ["O"], ["05361"], ["06"], ["O"], ["060536"], ["00"]),
    ("449", "80222", ["O", "O"], ["01718", "05492"], ["06", "06"], ["O", "O"], ["060171", "060549"], ["00", "00"]),
    ("451", "00005", ["G", "O"], ["5568", "04351"], ["10", "09"], ["G", "O"], ["105568", "090435"], ["09", "00"])
]
schema = "`WB-API-CNTY` string, `WB-API-UNIQUE` string, `WB-OIL-CODE` array<string>, `WB-OIL-LSE-NBR` array<string>, `WB-OIL-DIST` array<string>, `WB-GAS-CODE` array<string>, `WB-GAS-RRC-ID` array<string>, `WB-GAS-DIS` array<string>"
df = spark.createDataFrame(source_data, schema=schema)

# 核心处理流程
result_df = df.withColumn(
    "tmp_struct_arr",
    F.explode_outer(
        F.arrays_zip(
            "WB-OIL-CODE",
            "WB-OIL-LSE-NBR",
            "WB-OIL-DIST",
            "WB-GAS-CODE",
            "WB-GAS-RRC-ID",
            "WB-GAS-DIS"
        )
    )
).select(
    "WB-API-CNTY",
    "WB-API-UNIQUE",
    F.col("tmp_struct_arr.WB-OIL-CODE").alias("WB-OIL-CODE"),
    F.col("tmp_struct_arr.WB-OIL-LSE-NBR").alias("WB-OIL-LSE-NBR"),
    F.col("tmp_struct_arr.WB-OIL-DIST").alias("WB-OIL-DIST"),
    F.col("tmp_struct_arr.WB-GAS-CODE").alias("WB-GAS-CODE"),
    F.col("tmp_struct_arr.WB-GAS-RRC-ID").alias("WB-GAS-RRC-ID"),
    F.col("tmp_struct_arr.WB-GAS-DIS").alias("WB-GAS-DIS")
)

# 输出结果校验
result_df.show()

注意事项

  1. arrays_zip默认按最短数组长度截断打包,和Python原生zip逻辑一致,如果业务要求所有数组列长度必须相同,可以在处理前增加校验:通过F.size()计算各数组列长度,过滤长度不一致的异常数据或抛出告警。
  2. 必须使用explode_outer,如果用普通explode,空数组的行会被直接过滤,不符合空值保留的要求。
  3. 该方案时间复杂度为O(n),性能远高于多列单独explode后join的实现,不会出现数据错位问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 14:36:22