PySpark加载分区Parquet数据性能优化求助
性能优化方案
1. 批量读取所有目标分区路径(推荐)
直接生成日期范围内所有分区的路径列表,一次性交给Spark读取,彻底避免循环unionAll带来的执行计划冗余和多次IO调度开销:
import pandas as pd start_date = inputdate end_date = inputend # 生成所有日期对应的分区路径 path_list = [] for single_date in pd.date_range(start_date, end_date, freq='D'): yy = single_date.strftime("%y") mm = single_date.strftime("%m") dd = single_date.strftime("%d") path = f"abfs://XXXXX.dfs.core.windows.net/coredb/user_action/{yy}/{mm}/{dd}/user_action.parquet" path_list.append(path) # 一次性读取所有路径,自动合并为单个DataFrame df_union = spark.read.parquet(*path_list) \ .select("user_action_id", "account_id", "inserted" , "partner_id", "status", "service_id")
2. 利用路径通配符简化代码
如果日期范围是连续的月份/日期,可以用通配符减少路径生成的步骤,进一步提升代码简洁度:
- 同一月份内的日期范围:
df_union = spark.read.parquet("abfs://XXXXX.dfs.core.windows.net/coredb/user_action/23/01/{01..10}/user_action.parquet") \ .select("user_action_id", "account_id", "inserted" , "partner_id", "status", "service_id")
- 跨月份的日期范围:
df_union = spark.read.parquet( "abfs://XXXXX.dfs.core.windows.net/coredb/user_action/23/01/{25..31}/user_action.parquet", "abfs://XXXXX.dfs.core.windows.net/coredb/user_action/23/02/{01..10}/user_action.parquet" ).select("user_action_id", "account_id", "inserted" , "partner_id", "status", "service_id")
3. 启用分区发现(适用于Spark分区表)
如果数据是通过Spark partitionBy("yy", "mm", "dd")写入的分区表,可以直接读取根目录,通过过滤分区字段筛选日期范围,Spark会自动跳过无关分区:
# 读取分区表根目录 df = spark.read.parquet("abfs://XXXXX.dfs.core.windows.net/coredb/user_action/") # 转换日期为分区字段格式 start_part = start_date.strftime("%y%m%d") end_part = end_date.strftime("%y%m%d") # 过滤目标日期范围的分区 df_union = df.filter( (concat(df.yy, df.mm, df.dd) >= start_part) & (concat(df.yy, df.mm, df.dd) <= end_part) ).select("user_action_id", "account_id", "inserted" , "partner_id", "status", "service_id")
优化核心逻辑
循环unionAll会让Spark每次迭代都生成新的执行计划,且多次小文件读取会大幅增加IO调度成本;而批量读取、通配符、分区发现都是利用Spark的分区感知能力,一次性规划所有数据的读取任务,并行处理多个文件,从根本上降低开销、提升效率。
内容的提问来源于stack exchange,提问作者RobbeVL
相关产品推荐
相关产品推荐

