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

Pyspark无循环实现Array嵌套列展开为两列结构化DataFrame

PySpark 无循环展开嵌套数组列实现

实现思路

全程使用Spark原生内置算子实现,不依赖Python侧for/while循环、不使用带循环逻辑的自定义函数,所有计算由Spark Catalyst引擎优化为分布式执行逻辑,性能远高于循环类实现。
核心用到两个原生能力:

  • explode 函数:将数组类型列的每个外层元素拆分为独立行,是Spark原生行转列算子,无用户侧显式循环逻辑
  • 原生数组索引取值:直接通过下标读取二元数组的两个元素,不需要遍历数组内容

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, col

# 初始化Spark会话
spark = SparkSession.builder.master("local[*]").appName("expand_nested_array").getOrCreate()

# 构造输入样例数据
source_data = [
    ("a", ["e", "f", "g"], [["e", "f"], ["e", "g"], ["f", "g"]]),
    ("b", ["e", "f", "g", "h"], [["e", "f"], ["e", "g"], ["f", "g"], ["f", "h"]]),
    ("c", ["b", "c"], [["b", "c"]])
]
source_df = spark.createDataFrame(source_data, schema=["number", "values", "combination"])

# 核心转换逻辑
result_df = source_df.select(
    # 第一步:炸开嵌套数组的外层,每个二元组合成为单独一行
    explode(col("combination")).alias("pair")
).select(
    # 第二步:按下标取出二元组合的两个元素,映射为目标字段
    col("pair")[0].alias("value1"),
    col("pair")[1].alias("value2")
)

# 打印结果验证
result_df.show(truncate=False)

运行输出

执行后返回结果完全匹配预期,共8行数据:

+------+------+
|value1|value2|
+------+------+
|e     |f     |
|e     |g     |
|f     |g     |
|e     |f     |
|e     |g     |
|f     |g     |
|f     |h     |
|b     |c     |
+------+------+

合规说明

  • 代码全程未使用任何循环类函数,所有操作均为Spark原生分布式算子
  • 不需要将数据拉取到Driver端处理,可支持超大规模数据集的转换
  • 输出字段、行数、内容顺序完全匹配需求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 05:36:25