如何使用Polars将存储中的Parquet文件重分组为指定大小DataFrame
使用Polars将DataFrame重新拆分为特定大小的分区
你需要将多个Parquet文件合并后拆分为指定大小的DataFrame,作为临时数据仓库供Power BI和Python访问,原代码中的split_partitions伪方法可通过以下逻辑实现:
from pathlib import Path import uuid import polars as pl def repartition(directory_to_repartition, target_size_mb): repart_dir = Path(directory_to_repartition) old_paths = [v for v in repart_dir.iterdir() if v.suffix == '.parquet'] frames = [pl.read_parquet(path) for path in old_paths] big_frame = pl.concat(frames) # 假设内存足够容纳合并后的DataFrame # 按目标大小拆分DataFrame total_size_mb = big_frame.estimated_size(unit="mb") # 计算分区数,确保至少保留1个分区 num_partitions = max(1, int(total_size_mb / target_size_mb)) rows_per_partition = big_frame.height // num_partitions new_frames = [] for i in range(num_partitions): start_idx = i * rows_per_partition # 最后一个分区包含剩余所有行,避免过小分区 end_idx = (i + 1) * rows_per_partition if i != num_partitions - 1 else big_frame.height partition = big_frame.slice(start_idx, end_idx - start_idx) new_frames.append(partition) # 写入新的Parquet文件 for frame in new_frames: frame.write_parquet(repart_dir / f"{uuid.uuid4()}.parquet") # 清理旧文件 for old in old_paths: try: old.unlink() except FileNotFoundError: pass
关键说明:
target_size_mb参数以MB为单位,基于DataFrame的内存占用计算分区数,Parquet写入后因压缩会略小于该值;若需精确匹配文件大小,可增加临时文件写入校验逻辑(会额外消耗IO资源)。- 使用
estimated_size快速获取内存占用,避免全量数据扫描。 - 采用行数均分的拆分逻辑,最后一个分区自动兜底剩余行,保证分区大小相对均匀。
内容的提问来源于stack exchange,提问作者Isaacnfairplay
相关产品推荐
相关产品推荐

