PySpark 3.5.0计算DataFrame大小报错,求分区大小最优计算方案
PySpark 3.5.0 计算DataFrame大小与Parquet分区数确定方案
你使用的旧版估算代码依赖Spark私有API,在3.5.0版本中因内部API变更报错,私有API无兼容性保障,不建议使用。以下是官方支持或稳定的计算方法,以及Parquet分区数的确定建议:
一、DataFrame大小计算方法
1. 基于DESCRIBE EXTENDED的官方统计方法
这是Spark官方支持的方式,无需依赖私有API:
from pyspark.sql import DataFrame def get_df_size(input_df: DataFrame) -> int: input_df.createOrReplaceTempView("temp_stats_df") stats_result = input_df.sparkSession.sql("DESCRIBE EXTENDED temp_stats_df") size_row = stats_result.filter(stats_result.col_name == "sizeInBytes").first() return int(size_row["data_type"]) if size_row else 0
- 注意:若未执行过
ANALYZE TABLE temp_stats_df COMPUTE STATISTICS,返回的是Spark估算值;执行统计命令后可获得更精确的实际大小。
2. 缓存后获取实际内存占用
适合需要精确内存大小的场景(小到中等规模DataFrame):
def get_cached_df_size(input_df: DataFrame) -> int: cached_df = input_df.cache() cached_df.count() # 触发缓存加载 mem_size = 0 if cached_df.storageLevel.useMemory: storage_status = cached_df._jdf.storageStatus() mem_size = sum(status.memSize() for status in storage_status) # 可选:用完后释放缓存 cached_df.unpersist() return mem_size
- 该方法会将DataFrame加载到内存,大表使用需注意集群内存容量。
3. 抽样估算单分区大小
若仅需确定分区数,可通过抽样计算单分区平均大小:
def estimate_avg_partition_size(input_df: DataFrame, sample_count: int = 3) -> float: # 取前N个分区抽样 sample_df = input_df.limit(input_df.rdd.getNumPartitions() // sample_count) sample_size = get_cached_df_size(sample_df) return sample_size / sample_count
二、Parquet写入分区数建议
Parquet文件的最优大小通常为128MB~256MB(匹配主流存储系统的块/对象大小),可通过以下公式计算合适的分区数:
def get_optimal_partitions(df_size_bytes: int, target_mb_per_partition: int = 128) -> int: target_bytes = target_mb_per_partition * 1024 * 1024 return max(1, (df_size_bytes + target_bytes - 1) // target_bytes) # 向上取整
例如:10GB的DataFrame,按128MB分区,需10*1024/128=80个分区。
写入时可结合repartition或repartitionByRange调整,同时优先按业务分区键拆分,再按大小调整,避免生成大量小文件或超大文件。
内容的提问来源于stack exchange,提问作者OnTheWheels
相关产品推荐
相关产品推荐

