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

Python中如何实现超内存大体积CSV流式转换为Parquet文件

流式读取CSV并写入Parquet的实现方案

方案1:使用原生pyarrow实现(无需额外依赖)

你提到的Parquet流式写入能力pyarrow实际已经提供,核心用到pyarrow.parquet.ParquetWriter类,它支持逐批次写入RecordBatch或者小的Table,完美适配CSVStreamingReader的输出,无需把全量数据加载进内存。

import pyarrow as pa
import pyarrow.csv as pv
import pyarrow.parquet as pq

def stream_csv_to_parquet(csv_path, parquet_path, read_block_size=1 << 24):
    # 配置CSV读取参数,block_size控制每次读入内存的块大小,默认16MB可按需调整
    read_options = pv.ReadOptions(block_size=read_block_size)
    # 可选:提前定义CSV schema,避免流式读取时类型推断错误
    # csv_schema = pa.schema([
    #     ("col1", pa.int64()),
    #     ("col2", pa.string()),
    #     其他字段定义
    # ])
    # parse_options = pv.ParseOptions()
    # convert_options = pv.ConvertOptions(schema=csv_schema)
    
    # 初始化流式CSV读取器
    with pv.open_csv(csv_path, read_options=read_options) as reader:
        # 初始化Parquet写入器,用读取器的schema定义Parquet的schema
        with pq.ParquetWriter(parquet_path, schema=reader.schema) as writer:
            # 逐批次读取并写入
            for batch in reader:
                writer.write_batch(batch)
  • 注意点:
    • 如果CSV字段类型复杂容易推断错误,可以提前手动定义schema传入ConvertOptions,避免不同批次读取时类型不一致报错
    • block_size参数可以根据你的内存大小调整,越大批次越少、写入效率越高,占用内存也越大

方案2:使用polars实现(代码更简洁,性能更优)

如果可以接受引入额外库,polars对大文件流式处理的封装更易用,代码量更少:

import polars as pl

def stream_csv_to_parquet(csv_path, parquet_path):
    # 流式读取CSV,pl.scan_csv返回延迟执行的LazyFrame,不会加载全量数据
    df = pl.scan_csv(csv_path)
    # 直接流式写入Parquet,自动分批次处理,无需手动操作分片
    df.sink_parquet(parquet_path)
  • 优势:
    • 底层同样基于arrow实现,性能比手动写pyarrow逻辑更高
    • 自动处理schema推断、类型对齐问题,不用手动处理各种配置
    • 支持中途添加数据转换逻辑,比如过滤、字段修改等操作,同样全程流式不占内存

补充说明:你之前找到的RecordBatchFileWriter是用来写Arrow原生格式的,不要和Parquet的写入器搞混,上面两种方案都可以完美处理远大于内存的CSV文件转Parquet的需求。

内容的提问来源于stack exchange,提问作者OneRaynyDay

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 11:48:03