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()
注意事项
arrays_zip默认按最短数组长度截断打包,和Python原生zip逻辑一致,如果业务要求所有数组列长度必须相同,可以在处理前增加校验:通过F.size()计算各数组列长度,过滤长度不一致的异常数据或抛出告警。- 必须使用
explode_outer,如果用普通explode,空数组的行会被直接过滤,不符合空值保留的要求。- 该方案时间复杂度为O(n),性能远高于多列单独explode后join的实现,不会出现数据错位问题。
内容的提问来源于stack exchange,提问作者Eric
相关产品推荐
相关产品推荐

