Python 处理超大 CSV/Parquet 文件的内存管理:分块读取与流式写入

封面信息图

在数据分析、机器学习特征工程与 ETL 数据清洗中,Python 的 Pandas 库无疑是大家最熟悉的工具。

但在处理超大规模数据集(例如一个 20GB 的历史日志 CSV,或者 50GB 的 Parquet 数据集)时,很多初中级工程师最容易写出的一行代码就是:

# 致命代码:在 16GB 内存的服务器上直接全量载入 20GB 文件
df = pd.read_csv("huge_dataset_20gb.csv")

其结果必然是:操作系统内存瞬间被吃光,Swap 分区爆满,最终触发内核 Out of Memory: Kill process (OOM),整个进程被系统粗暴杀死!

很多人的误区在于以为“20GB 的文件载入内存只要 20GB”:
实际上,Pandas 在将字符串、对象转换为内部 DataFrame 结构时,存在大量的对象封装开销,内存膨胀系数通常在 3 到 5 倍!一个 20GB 的 CSV 文件在 Pandas 内存中往往需要消耗 60GB ~ 100GB 的物理内存!

今天我们拆解如何利用 Pandas 分块迭代(Chunking)、PyArrow 流式写入与 Polars 惰性流式引擎(Lazy Streaming),在仅有 2GB 内存的小服务器上,从容流式处理 100GB 级别的超大文件。


一、全量载入 vs 流式处理的内存占用对比

flowchart TD
    subgraph Full_Load [全量载入内存 (容易 OOM 崩溃)]
        F1[20GB 磁盘大文件] -->|一次性全部载入| F2[内存膨胀至 80GB -> 触发内核 OOMKilled!]
    end
    
    subgraph Stream_Chunk [分块流式迭代 (内存恒定在 200MB!)]
        S1[20GB 磁盘大文件] --> S2[迭代器读取 Chunk 1 (10万行, 150MB)]
        S2 --> S3[清洗与特征计算]
        S3 --> S4[追加写入目标 Parquet]
        S4 --> S5[释放 Chunk 1 内存 -> 读取 Chunk 2]
    end

二、方案 1:Pandas 基于 chunksize 的优雅分块清洗实战

import pandas as pd
from typing import Callable

def process_huge_csv_by_chunks(
    input_csv_path: str,
    output_parquet_path: str,
    chunk_rows: int = 100000,
    transform_func: Callable[[pd.DataFrame], pd.DataFrame] = None
):
    print(f"[*] 开始以每块 {chunk_rows} 行分块流式处理: {input_csv_path}")
    
    # 1. 核心参数:chunksize 返回一个可迭代的 TextFileReader 生成器对象
    csv_reader = pd.read_csv(
        input_csv_path,
        chunksize=chunk_rows,
        low_memory=True,
        dtype={
            "user_id": "int32",       # 显式压缩数据类型(int64 -> int32 节省 50% 内存)
            "is_vip": "bool",
            "amount_cents": "int32"
        }
    )

    is_first_chunk = True
    total_processed = 0

    for idx, chunk_df in enumerate(csv_reader, start=1):
        # 执行业务清洗与特征工程
        if transform_func:
            chunk_df = transform_func(chunk_df)

        # 2. 流式追加写入 Parquet 或 CSV
        # 若写入 CSV: 使用 mode='a', header=is_first_chunk
        chunk_df.to_parquet(
            f"/tmp/chunk_{idx}.parquet",
            engine="pyarrow",
            compression="snappy"
        )

        total_processed += len(chunk_df)
        print(f"[✓] 已完成第 {idx} 块处理,累计处理 {total_processed} 行,当前常驻内存恒定 < 300MB。")

    print(f"[Finished] 全量 {total_processed} 行数据处理大功告成!")

三、方案 2:利用 Polars 惰性流式计算(Lazy Streaming Engine)

在 2026 年的现代数据工程中,更推荐使用基于 Rust 构建的 Polars。其内置的 scan_csv + sink_parquet 能够在流式模式(Streaming Mode)下自动完成多核并发分块与流水线调度,比 Pandas 快 5 到 10 倍:

import polars as pl

def process_huge_file_with_polars_streaming(input_csv: str, output_parquet: str):
    print("[*] 启动 Polars 惰性流式计算引擎...")
    
    # 1. 扫描元数据建立计算图 (零数据载入内存)
    lazy_query = (
        pl.scan_csv(input_csv)
        .filter(pl.col("amount_cents") > 0)
        .with_columns([
            (pl.col("amount_cents") / 100.0).alias("amount_yuan"),
            pl.col("user_id").cast(pl.Int32)
        ])
    )

    # 2. 核心:开启 streaming=True 流式下沉写入磁盘!
    # Polars 会自动在底层以微批次流式抽取并利用多核 SIMD 指令加速,内存占用极低!
    lazy_query.sink_parquet(
        output_parquet,
        compression="zstd",
        maintain_order=False
    )
    print(f"[✓] Polars 流式处理完毕,输出已保存至 {output_parquet}")

四、生产级内存压榨核心黄金法则

  1. 强行下沉数据类型(Downcasting Types)
    • 将默认的 float64 / int64 下沉为 float32 / int32,直接立省 50% 内存
    • 对低基数字符串列(如性别、省份、状态码)使用 category 分类类型,可将字符串内存占用压缩 80% 以上
  2. 选择列裁剪读取(usecols:如果表有 80 列,业务只用其中 5 列,务必配置 usecols=["id", "amount", "dt"],杜绝载入 75 列无用脏数据;
  3. 主动调用 gc.collect():在处理完每个超大 DataFrame 块并将其删除(del chunk_df)后,显式触发 Python 垃圾回收,加速内存交还给操作系统。

掌握了分块迭代与流式下沉技术,哪怕在最小配置的轻量云服务器上,你也能举重若轻地驾驭上百 GB 级别的大数据清洗流水线。

Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐