如何不将全量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
相关产品推荐
相关产品推荐

