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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 21:01:53