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

如何高效构建含Hive风格分区最新数据的PySpark DataFrame

从Hive风格分区Parquet中提取每个source_id的最新year_month数据

数据存储结构

数据采用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数据:

  1. 生成仅包含source_id和year_month的DataFrame,按source_id分组取year_month的最大值;
  2. 将结果与全量数据关联,筛选目标数据。

实现代码如下:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 14:15:56