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参数可以根据你的内存大小调整,越大批次越少、写入效率越高,占用内存也越大
- 如果CSV字段类型复杂容易推断错误,可以提前手动定义schema传入
方案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
相关产品推荐
相关产品推荐

