如何基于df_grouped动态Union指定分区的PySpark DataFrame
问题
我的数据湖文件按partition_Continent和partition_Country两个分区存储,现有df_grouped提供筛选所需的分区条件(示例返回2条记录)。当前我通过生成filter_condition,用where子句读取对应分区的数据,但认为直接读取指定分区再执行Union的效率更高(手动指定分区路径后Union的示例代码如下)。请问如何基于df_grouped的返回结果实现这种动态Union操作?
当前实现filter_condition的代码:
filter_condition = " OR ".join( [ ( f"(partition_Continent = '{i.Continent}'" f" AND partition_Country = '{i.Country}')" ) for i in df_grouped.distinct().collect() ] )
当前读取数据的代码:
df_presented = spark.read.parquet(f'abfss://laketommy@tommy.dfs.core.windows.net/Tommy').where(filter_condition)
目标手动示例代码:
df_presented = spark.read.parquet(f'abfss://laketommy@tommy.dfs.core.windows.net/Tommy/partition_Continent=Europe/partition_Country=UK').union(spark.read.parquet(f'abfss://laketommy@tommy.dfs.core.windows.net/Tommy/partition_Continent=Asia/partition_Country=China'))
解决方案
你可以通过遍历df_grouped的分区记录,动态生成每个分区的路径,逐个读取DataFrame后再执行批量Union操作,具体实现如下:
1. 获取去重后的分区键值对
先从df_grouped中提取去重的Continent和Country组合,避免重复读取同一分区:
# 收集去重后的分区条件 partition_list = df_grouped.select("Continent", "Country").distinct().collect()
2. 动态生成分区路径并读取DataFrame
遍历分区列表,为每个组合拼接完整的分区路径,读取对应的Parquet文件,将所有DataFrame存入列表:
base_path = "abfss://laketommy@tommy.dfs.core.windows.net/Tommy" df_list = [] for item in partition_list: # 拼接分区路径 partition_path = f"{base_path}/partition_Continent={item.Continent}/partition_Country={item.Country}" # 读取该分区的DataFrame df_part = spark.read.parquet(partition_path) df_list.append(df_part)
3. 批量Union所有DataFrame
利用functools.reduce对列表中的DataFrame执行批量Union操作:
from functools import reduce from pyspark.sql import DataFrame # 批量合并所有分区的DataFrame df_presented = reduce(DataFrame.union, df_list)
补充说明
- 这种方式直接读取指定分区路径,Spark可以跳过全局分区元数据扫描,在分区数量较多时,效率明显高于
where过滤的方式。 - 如果你的Spark版本在2.3及以上,可使用
unionByName替代union,确保列顺序不一致时也能正确合并(需调用DataFrame.unionByName)。
内容的提问来源于stack exchange,提问作者tommyhmt
相关产品推荐
相关产品推荐

