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

如何不将全量csv加载到内存即可将其转换为.arrow格式文件

问题描述

我找到过一个类似的问题,里面提到了sys.stdin.buffer和sys.stdout.buffer,但我不清楚怎么用这两个方法写入并命名.arrow文件。我在PyArrow官方文档里也没找到对应的说明。

我要处理的CSV文件没有空值,但带有时间戳索引,文件大小约100GB,完全没办法加载到内存中处理。我尝试写了代码处理,但每次循环都会覆盖之前生成的文件。

这是我第一次发帖,非常感谢所有贡献者,我之前遇到的99.9%的问题都已经被大家解答过了。

我的代码:

import sys

import pandas as pd
import pyarrow as pa

SPLIT_ROWS = 1     ### 测试用小分块,每次读1行

def main():
    writer = None
    for split in pd.read_csv(sys.stdin.buffer, chunksize=SPLIT_ROWS):

        table = pa.Table.from_pandas(split)
        # 写入文件
        with pa.OSFile('test.arrow', 'wb') as sink:     ### 目前没有追加模式
            with pa.RecordBatchFileWriter(sink, table.schema) as writer:
                writer.write_table(table)
    writer.close()

if __name__ == "__main__":
    main()

命令行调用方式:

cat data.csv | python test.py

解决方案

代码问题核心是把文件句柄和Arrow写入器的初始化放在了分块遍历的循环内部,每次循环都会以覆盖模式重新创建test.arrow,之前写入的内容自然会被清空。

针对100GB大文件的场景,提供两种可行方案:

方案1:输出单个完整Arrow文件

如果需要最终输出一个完整的Arrow文件,只要把写入器的初始化移到循环外即可:

import sys
import pandas as pd
import pyarrow as pa

# 可根据内存容量调整分块大小,100GB文件建议设为100万~500万行
SPLIT_ROWS = 1000000 

def main():
    writer = None
    sink = pa.OSFile('test.arrow', 'wb')
    try:
        for split in pd.read_csv(sys.stdin.buffer, chunksize=SPLIT_ROWS):
            table = pa.Table.from_pandas(split)
            # 首次读取到分块时,基于表结构初始化写入器
            if writer is None:
                writer = pa.RecordBatchFileWriter(sink, table.schema)
            writer.write_table(table)
    finally:
        if writer is not None:
            writer.close()
        sink.close()

if __name__ == "__main__":
    main()

方案2:输出多个分块Arrow文件

如果后续需要分布式处理,希望生成多个独立的分块Arrow文件,给每个分块的文件名加上序号即可:

import sys
import pandas as pd
import pyarrow as pa

SPLIT_ROWS = 1000000 

def main():
    chunk_index = 0
    for split in pd.read_csv(sys.stdin.buffer, chunksize=SPLIT_ROWS):
        table = pa.Table.from_pandas(split)
        # 每个分块生成独立的文件名
        file_path = f'test_part_{chunk_index}.arrow'
        with pa.OSFile(file_path, 'wb') as sink:
            with pa.RecordBatchFileWriter(sink, table.schema) as writer:
                writer.write_table(table)
        chunk_index += 1

if __name__ == "__main__":
    main()

两种方案的命令行调用方式都不需要改动,和原有用法一致。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 18:36:04