资讯动态

如何用 PyArrow 的 IPC 接口写入和读取流式与文件两种 Arrow 序列化格式

发布时间:2026/9/14 3:46:19 来源:尧图企业网站定制
如何用 PyArrow 的 IPC 接口写入和读取流式与文件两种 Arrow 序列化格式【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow当你需要把 PyArrow 的 record batch 序列化成二进制数据——存到磁盘、写到内存缓冲、或者通过网络发给另一端时——PyArrow 的 IPC 接口提供了两条现成的路径pyarrow.ipc.new_stream创建流式格式pyarrow.ipc.new_file创建文件随机访问格式。本文按照 docs/source/python/ipc.rst 的实际演示走一遍构造 batch → 写出 → 读回 → 验证一致的完整过程并在文档覆盖的范围内说明两种格式的差异和读取大文件时的内存映射技巧。前提条件已按 docs/source/python/install.rst 安装好 pyarrowWindows、Linux、macOS 下可用pip install pyarrow。如果后面要用read_pandas安装文档列出的可选依赖要求 pandas 2.2.2 或更高版本。文档同时建议先阅读 docs/source/python/memory.rst 的 Memory and IO 部分因为写目标sink都是其中的 IO 对象。先选格式流式还是文件Arrow 定义了两种序列化 record batch 的二进制格式摘自 ipc 文档Streaming format流式格式用于发送任意长度的 record batch 序列。必须从头到尾按顺序处理不支持随机访问File or Random Access format文件/随机访问格式用于序列化固定数量的 record batch。支持随机访问配合 memory map 使用时很有用。判断依据是读端需求如果数据是顺序消费的比如写入 socket、逐批转发用流式如果读端需要知道总批数、直接取第 N 个 batch或文件会落在磁盘上配合 memory map 使用用文件格式。本文两条路径都会走一遍。构造待序列化的 record batch两种格式的写者 API 相同示例从同一个 batch 出发import pyarrow as pa data [ pa.array([1, 2, 3, 4]), pa.array([foo, bar, baz, None]), pa.array([True, None, False, True]) ] batch pa.record_batch(data, names[f0, f1, f2]) batch.num_rows # 4 batch.num_columns # 3这个 batch 有 3 列int64、string、bool、4 行后面无论写多少次读回后都可以和它做逐字节比较。写入并读取流式格式流式写法创建pa.BufferOutputStream作为内存 sink文档指出换成 socket 等任意可写目标同样成立用pa.ipc.new_stream打开写者连续写入 5 个 batchsink pa.BufferOutputStream() with pa.ipc.new_stream(sink, batch.schema) as writer: for i in range(5): writer.write_batch(batch) buf sink.getvalue() # 完整流内容内存中的字节 Buffer注意两点都来自文档说明创建StreamWriter时必须传入 schema因为同一股流中所有 batch 的 schema列名和类型必须一致文档中示例写入 5 个这样的 batch 后buf.size为 1984 字节文档示例值与具体平台/版本有关不要当作固定预期。读回用pa.ipc.open_stream或pyarrow.RecordBatchStreamReaderwith pa.ipc.open_stream(buf) as reader: schema reader.schema batches [b for b in reader] schema # f0: int64 # f1: string # f2: bool len(batches) # 5 batches[0].equals(batch) # True验证方式就是文档使用的equals读回的每个 batch 与原始输入完全相等即说明序列化往返无损。另外文档指出一个对性能有直接影响的行为如果输入源支持零拷贝读取如 memory map 或pyarrow.BufferReader读回的 batch 也是零拷贝的读取时不分配任何新内存。写入并读取文件随机访问格式文件写者pa.ipc.new_file与流式写者 API 相同这里示例写入 10 个 batchsink pa.BufferOutputStream() with pa.ipc.new_file(sink, batch.schema) as writer: for i in range(10): writer.write_batch(batch) buf sink.getvalue()文档中示例buf.size为 4226 字节同样是文档示例值。读端的区别在于RecordBatchFileReader要求输入源必须有seek方法以支持随机访问而流式读者只需要读操作。由于能访问完整 payload文件读者可以直接给出批数并随机取任意一个with pa.ipc.open_file(buf) as reader: num_record_batches reader.num_record_batches b reader.get_batch(3) num_record_batches # 10 b.equals(batch) # True验证口径与流式一致num_record_batches等于写入次数10get_batch(3)取到的第 4 个 batch 与原始 batchequals为True。可选直接读成 pandas DataFrame两种格式的文件/流读者都提供read_pandas方法把多个 record batch 读入并合并成单个 DataFramewith pa.ipc.open_file(buf) as reader: df reader.read_pandas() df[:5]文档给出的示例输出标注为示例结果列内容随写入数据变化f0 f1 f2 0 1 foo True 1 2 bar None 2 3 baz False 3 4 NaN True 4 1 foo True写入磁盘大文件并配合 memory map 读取当数据大到不能一次装进内存时ipc 文档给出的做法是写端按 batch 分块示例为 1000 个 batch、每批 10000 个 int32共 10M 整数sink 换成pa.OSFile直接写磁盘读端再选择普通文件读或 memory map 读。BATCH_SIZE 10000 NUM_BATCHES 1000 schema pa.schema([pa.field(nums, pa.int32())]) with pa.OSFile(bigfile.arrow, wb) as sink: with pa.ipc.new_file(sink, schema) as writer: for row in range(NUM_BATCHES): batch pa.record_batch([pa.array(range(BATCH_SIZE), typepa.int32())], schema) writer.write(batch)记录 batch 支持多列实践中写的是等价的 Table。分块写的好处是写端理论上只需把当前 batch 留在内存里。读回有两种方式。普通文件方式每次读取会分配新内存类似 Python 文件对象文档示例输出with pa.OSFile(bigfile.arrow, rb) as source: loaded_array pa.ipc.open_file(source).read_all() print(LEN:, len(loaded_array)) # LEN: 10000000 print(RSS: {}MB.format(pa.total_allocated_bytes() 20)) # RSS: 38MB以上为文档示例输出RSS 数值取决于环境与运行时状态不要当作固定成功标准LEN应等于写入的总行数 10000000。改用pa.memory_map后Arrow 直接引用映射自磁盘的数据不为自己分配内存操作系统按需换入页面、内存压力大时无写回成本地换出从而更容易读入超过总内存的数组with pa.memory_map(bigfile.arrow, rb) as source: loaded_array pa.ipc.open_file(source).read_all() print(LEN:, len(loaded_array)) # LEN: 10000000 print(RSS: {}MB.format(pa.total_allocated_bytes() 20)) # RSS: 0MBRSS: 0MB 同样是文档示例输出表示该次运行中total_allocated_bytes()未计入映射内存。文档末尾附带一条边界说明pyarrow.parquet.read_table等高层 API 也提供memory_map选项但那种情况下 memory mapping 不能帮助降低常驻内存消耗详情见文档中引用的parquet_mmap一节。两种格式的核对清单项目流式new_stream/open_stream文件new_file/open_filebatch 数量任意长度的序列固定数量读端可通过num_record_batches得知随机访问不支持只能从头到尾顺序处理支持get_batch(i)取任意 batch输入源要求只需要读操作必须有seek方法典型场景socket 等顺序消费落盘 memory map验证方式batches[i].equals(batch)get_batch(3).equals(batch)两条路径的验证结论都落在equals返回True流式往返 5 个 batch、文件随机访问取第 4 个 batch均与原始输入相等。如果后续要处理带压缩或网络传输的 IPC 场景可以在 docs/source/python/memory.rst 中查看CompressedInputStream/CompressedOutputStream等 IO 对象它们是文档列出的NativeFile家族成员可作为本文 sink/source 的替换目标。【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价