在PySpark的实际数据处理场景中,经常会遇到包含嵌套数组的复杂数据结构,比如JSON解析后的字段、业务系统导出的多层嵌套数据等。展开这些嵌套数组是提取有效信息的必要步骤,但如果操作方式不合理,很容易引发数据行数指数级增长的问题,也就是所谓的数据膨胀,这会直接导致后续的计算任务耗时剧增,甚至因为内存不足而失败。

嵌套数组展开导致数据膨胀的原因
PySpark中展开数组最常用的函数是explode,它的作用是将数组中的每个元素拆成单独的一行,同时保留其他列的值。当处理单层数组时,这个函数的表现是符合预期的,比如一个包含3个元素的数组展开后会生成3行数据。
但如果要展开多个嵌套数组,尤其是多个数组属于同一层级、元素之间存在对应关系时,直接多次调用explode就会出现问题。因为每次调用explode都会对当前的所有行进行展开,多个数组的展开会产生笛卡尔积效果,最终行数会是各个数组长度的乘积,当数组长度稍大时就会呈现指数级增长。
我们可以通过一个简单的示例来验证这个问题,首先构造测试数据:
from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, col, array_zip
# 初始化SparkSession
spark = SparkSession.builder
.appName("nested_array_expand")
.master("local[*]")
.getOrCreate()
# 构造测试数据,包含三个同长度的嵌套数组
test_data = [
(1, ["a", "b", "c"], [10, 20, 30], ["x", "y", "z"]),
(2, ["d", "e"], [40, 50], ["m", "n"])
]
# 定义表结构
schema = ["id", "arr1", "arr2", "arr3"]
# 创建DataFrame
df = spark.createDataFrame(test_data, schema=schema)
df.show(truncate=False)
上述代码创建的数据中,id为1的行包含三个长度都是3的数组,id为2的行包含三个长度都是2的数组。如果我们直接对这三个数组依次使用explode函数,看看结果行数:
# 错误示范:多次调用explode导致笛卡尔积
wrong_df = df.select(
col("id"),
explode(col("arr1")).alias("val1"),
explode(col("arr2")).alias("val2"),
explode(col("arr3")).alias("val3")
)
wrong_df.show(truncate=False)
print(f"错误展开后的总行数:{wrong_df.count()}")
执行后会发现,id为1的行最终生成了3*3*3=27行,id为2的行生成了2*2*2=8行,总行数远超过我们预期的3+2=5行,这就是典型的数据膨胀问题。
高效展开嵌套数组的几种方案
方案一:使用array_zip函数合并数组后展开
如果多个嵌套数组的元素是一一对应的关系,比如上述示例中arr1的第1个元素对应arr2的第1个元素、arr3的第1个元素,那么可以先用array_zip函数将多个数组合并成一个结构体数组,再对合并后的数组使用一次explode,这样就不会产生笛卡尔积。
实现代码如下:
# 正确方案1:使用array_zip合并数组后展开
from pyspark.sql.functions import array_zip, explode, col
correct_df1 = df.select(
col("id"),
# 将三个数组按位置合并为结构体数组
array_zip(col("arr1"), col("arr2"), col("arr3")).alias("zipped_arr")
)
# 展开合并后的数组
correct_df1 = correct_df1.select(
col("id"),
explode(col("zipped_arr")).alias("zipped_val")
)
# 提取结构体中的各个字段
correct_df1 = correct_df1.select(
col("id"),
col("zipped_val.arr1").alias("val1"),
col("zipped_val.arr2").alias("val2"),
col("zipped_val.arr3").alias("val3")
)
correct_df1.show(truncate=False)
print(f"方案1展开后的总行数:{correct_df1.count()}")
这个方案最终生成的行数是3+2=5行,符合我们的预期,没有产生数据膨胀。它的核心逻辑是将同位置的数组元素打包成一个结构体,只展开一次,避免了多次展开的笛卡尔积。
方案二:使用posexplode函数获取索引后关联
如果数组的对应关系不是严格的一一对应,或者需要更灵活的控制,可以使用posexplode函数,它和explode的区别是会同时返回数组元素的索引和值。我们可以对每个数组使用posexplode获取索引,然后通过索引进行关联,避免笛卡尔积。
实现代码如下:
from pyspark.sql.functions import posexplode, col
# 正确方案2:使用posexplode获取索引后关联
# 展开第一个数组,获取索引和值
step1_df = df.select(
col("id"),
posexplode(col("arr1")).alias("idx1", "val1")
)
# 展开第二个数组,获取索引和值
step2_df = df.select(
col("id"),
posexplode(col("arr2")).alias("idx2", "val2")
)
# 展开第三个数组,获取索引和值
step3_df = df.select(
col("id"),
posexplode(col("arr3")).alias("idx3", "val3")
)
# 按照id和索引进行关联,只保留索引相同的行
correct_df2 = step1_df.join(step2_df, (step1_df.id == step2_df.id) & (step1_df.idx1 == step2_df.idx2))
.join(step3_df, (step1_df.id == step3_df.id) & (step1_df.idx1 == step3_df.idx3))
.select(step1_df.id, "val1", "val2", "val3")
correct_df2.show(truncate=False)
print(f"方案2展开后的总行数:{correct_df2.count()}")
这个方案同样可以得到5行结果,适合数组长度不完全一致、需要按索引匹配的场景,不过多步join的操作会比方案一的性能稍差一些。
方案三:使用高阶函数inline展开
PySpark支持高阶函数,其中inline函数可以直接将结构体数组展开成多行,结合transform函数可以实现数组的预处理后再展开,适合更复杂的嵌套场景。
实现代码如下:
from pyspark.sql.functions import expr
# 正确方案3:使用高阶函数inline展开
correct_df3 = df.selectExpr(
"id",
"inline(array_zip(arr1, arr2, arr3))"
)
correct_df3.show(truncate=False)
print(f"方案3展开后的总行数:{correct_df3.count()}")
这个方案代码最简洁,array_zip合并数组后,inline直接将其展开,效果和方案一一致,而且执行效率更高,是推荐的使用方式。
不同方案的适用场景对比
我们可以通过表格对比几种方案的优缺点和适用场景:
| 方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| 多次explode | 代码简单,理解成本低 | 会产生笛卡尔积,导致数据膨胀 | 仅适用于展开单个数组的场景 |
| array_zip+explode | 避免数据膨胀,逻辑清晰 | 需要多步操作提取字段 | 多个同长度、一一对应的数组展开 |
| posexplode+join | 支持按索引灵活匹配 | 多步join性能较差,代码复杂 | 数组长度不一致、需要按索引匹配的场景 |
| inline+array_zip | 代码简洁,执行效率高 | 需要熟悉PySpark高阶函数语法 | 所有同位置数组展开的场景,优先推荐 |
避免数据膨胀的实用技巧
除了选择合适的展开方案,还有几个技巧可以帮助避免嵌套数组展开时的数据膨胀问题:
- 展开前先确认数组的对应关系:如果多个数组是同层级、一一对应的,绝对不要多次调用explode,优先使用array_zip合并后展开。
- 提前过滤无效数据:如果数组中包含大量无效元素,在展开前先使用
filter函数过滤掉,减少展开后的行数。 - 控制展开层级:如果嵌套层级过深,优先展开最外层的必要数组,避免一次性展开所有嵌套层级,分步骤处理可以降低数据膨胀的风险。
- 监控展开后的数据量:在开发阶段,可以先取少量数据(比如使用
limit函数)测试展开后的行数,确认符合预期后再处理全量数据。
总结
PySpark中嵌套数组展开的核心是避免多次explode带来的笛卡尔积问题,对于一一对应的多个数组,优先使用array_zip合并后配合inline或者explode展开,是最高效且不易出错的方式。在处理复杂嵌套结构时,先理清数组之间的对应关系,再选择合适的展开方案,就能有效避免指数级数据膨胀的问题,提升任务的执行效率。