如何高效构建含Hive风格分区最新数据的PySpark DataFrame
数据存储结构
数据采用Hive风格分区存储在Parquet文件中,分区键为source_id和year_month,目录结构如下:
/data/plandata/ ├── source_id=1/ │ ├── year_month=2023-08-01/ │ │ ├── file1.parquet │ │ └── ... │ ├── year_month=2023-09-01/ │ │ ├── file2.parquet │ │ └── ... ├── source_id=2/ │ ├── year_month=2023-09-01/ │ │ ├── file3.parquet │ │ └── ... │ ├── year_month=2023-10-01/ │ │ ├── file4.parquet │ │ └── ... └── ...
已实现的解决方案
通过以下步骤提取每个source_id对应的最新year_month数据:
- 生成仅包含
source_id和year_month的DataFrame,按source_id分组取year_month的最大值; - 将结果与全量数据关联,筛选目标数据。
实现代码如下:
from pyspark.sql import functions as F basepath = '.../test_data' # 生成测试数据(实际环境中数据已存在) data = [ {'source_id': 1, 'year_month': '2023-08-01', 'dataval': 'a-01-08'}, {'source_id': 1, 'year_month': '2023-08-01', 'dataval': 'b-01-08'}, {'source_id': 1, 'year_month': '2023-08-01', 'dataval': 'c-01-08'}, {'source_id': 1, 'year_month': '2023-09-01', 'dataval': 'a-01-09'}, {'source_id': 1, 'year_month': '2023-09-01', 'dataval': 'b-01-09'}, {'source_id': 1, 'year_month': '2023-09-01', 'dataval': 'c-01-09'}, {'source_id': 2, 'year_month': '2023-09-01', 'dataval': 'a-02-09'}, {'source_id': 2, 'year_month': '2023-09-01', 'dataval': 'b-02-09'}, {'source_id': 2, 'year_month': '2023-09-01', 'dataval': 'c-02-09'}, {'source_id': 2, 'year_month': '2023-10-01', 'dataval': 'a-02-10'}, {'source_id': 2, 'year_month': '2023-10-01', 'dataval': 'b-02-10'}, {'source_id': 2, 'year_month': '2023-10-01', 'dataval': 'c-02-10'}, {'source_id': 3, 'year_month': '2023-08-01', 'dataval': 'a-03-08'}, {'source_id': 3, 'year_month': '2023-08-01', 'dataval': 'b-03-08'}, {'source_id': 3, 'year_month': '2023-08-01', 'dataval': 'c-03-08'}, {'source_id': 3, 'year_month': '2023-10-01', 'dataval': 'a-03-10'}, {'source_id': 3, 'year_month': '2023-10-01', 'dataval': 'b-03-10'}, {'source_id': 3, 'year_month': '2023-10-01', 'dataval': 'c-03-10'}, ] df = spark.createDataFrame(data) df.write.mode('overwrite').partitionBy('source_id', 'year_month').parquet(basepath) # 提取最新数据的逻辑 df_recent_partitions = spark.read.parquet(basepath).select('source_id', 'year_month') df_recent_partitions = df_recent_partitions.groupBy('source_id' ).agg(F.max('year_month').alias('max_year_month') ) df_latest_data = spark.read.parquet(basepath) df_latest_data = df_latest_data.join(df_recent_partitions, ( (df_latest_data.source_id == df_recent_partitions.source_id) & (df_latest_data.year_month == df_recent_partitions.max_year_month) ) ).select(df_latest_data['*']) df_latest_data.show()
执行结果符合预期:
+-------+---------+----------+ |dataval|source_id|year_month| +-------+---------+----------+ |a-02-10| 2|2023-10-01| |b-02-10| 2|2023-10-01| |c-02-10| 2|2023-10-01| |a-03-10| 3|2023-10-01| |b-03-10| 3|2023-10-01| |c-03-10| 3|2023-10-01| |a-01-09| 1|2023-09-01| |b-01-09| 1|2023-09-01| |c-01-09| 1|2023-09-01| +-------+---------+----------+
疑问解答
1. 该方法是否高效?
假设a验证:正确
创建df_recent_partitions时,Spark确实不需要读取Parquet文件内容。因为source_id和year_month是分区列,Spark可以直接从目录名称中提取这两个字段,执行计划中的ReadSchema: struct<>也证明了这一点——没有读取文件内的任何数据,只扫描分区目录元数据。假设b验证:错误
当前的join写法无法让Spark自动跳过无关分区。因为join是在运行时进行的关联操作,Spark无法将join条件提前下推到分区过滤阶段。从执行计划可以看到,读取全量数据的FileScan中,PartitionFilters只有isnotnull(source_id#128), isnotnull(year_month#129),没有具体的(source_id, year_month)等值过滤,意味着Spark会读取所有分区的文件,再在内存中通过join筛选数据,这会浪费大量IO资源,效率较低。
优化方案:
先收集最新分区的条件,生成过滤表达式,直接读取对应分区的数据,避免全量扫描:
# 获取每个source_id的最新year_month,转换为列表 latest_partitions = df_recent_partitions.collect() # 生成过滤条件:(source_id=1 AND year_month='2023-09-01') OR (source_id=2 AND year_month='2023-10-01') ... filter_expr = F.lit(False) for row in latest_partitions: filter_expr = filter_expr | ((F.col("source_id") == row.source_id) & (F.col("year_month") == row.max_year_month)) # 直接过滤后读取数据 df_latest_data_optimized = spark.read.parquet(basepath).filter(filter_expr) df_latest_data_optimized.show()
这种方式下,Spark会将过滤条件下推到分区扫描阶段,只读取符合条件的分区文件,大幅提升效率。
2. 如何通过执行计划验证假设?
解读你提供的执行计划:
df_recent_partitions的执行计划:
== Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- HashAggregate(keys=[source_id#39], functions=[max(year_month#40)]) +- Exchange hashpartitioning(source_id#39, 200), ENSURE_REQUIREMENTS, [plan_id=449] +- HashAggregate(keys=[source_id#39], functions=[partial_max(year_month#40)]) +- FileScan parquet [source_id#39,year_month#40] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[...], PartitionFilters: [], PushedFilters: [], ReadSchema: struct<>关键看
FileScan部分:ReadSchema: struct<>表示没有读取Parquet文件的实际内容,仅从分区目录获取source_id和year_month,验证了假设a的正确性。PartitionFilters: []说明需要扫描所有分区来计算每个source_id的最大year_month,这是合理的。df_latest_data的执行计划:
== Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- Project [dataval#127, source_id#128, year_month#129] +- BroadcastHashJoin [source_id#128, year_month#129], [source_id#39, max_year_month#52], Inner, BuildRight, false :- FileScan parquet [dataval#127,source_id#128,year_month#129] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[...], PartitionFilters: [isnotnull(source_id#128), isnotnull(year_month#129)], PushedFilters: [], ReadSchema: struct<dataval:string> +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, int, true], input[1, date, false]),false), [plan_id=435] +- Filter isnotnull(max_year_month#52) +- HashAggregate(keys=[source_id#39], functions=[max(year_month#40)]) +- Exchange hashpartitioning(source_id#39, 200), ENSURE_REQUIREMENTS, [plan_id=431] +- HashAggregate(keys=[source_id#39], functions=[partial_max(year_month#40)]) +- FileScan parquet [source_id#39,year_month#40] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[...], PartitionFilters: [isnotnull(source_id#39)], PushedFilters: [], ReadSchema: struct<>读取全量数据的
FileScan中,PartitionFilters只有非空校验,没有具体的等值过滤条件,说明Spark没有将join条件转换为分区过滤,会读取所有分区的文件,验证了假设b不成立。
优化后的df_latest_data_optimized执行计划中,FileScan的PartitionFilters会包含具体的(source_id, year_month)等值条件,比如(source_id = 1 AND year_month = '2023-09-01') OR ...,证明Spark仅读取目标分区。
内容的提问来源于stack exchange,提问作者teejay

