如何用 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 安装好 pyarrow(Windows、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 个 batch:
sink = 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.RecordBatchStreamReader):
with 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 个 batch:
sink = 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等于写入次数(10),get_batch(3)取到的第 4 个 batch 与原始 batchequals为True。
可选:直接读成 pandas DataFrame
两种格式的文件/流读者都提供read_pandas方法,把多个 record batch 读入并合并成单个 DataFrame:
with 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), type=pa.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: 0MB(RSS: 0MB 同样是文档示例输出,表示该次运行中total_allocated_bytes()未计入映射内存。)
文档末尾附带一条边界说明:pyarrow.parquet.read_table等高层 API 也提供memory_map选项,但那种情况下 memory mapping 不能帮助降低常驻内存消耗,详情见文档中引用的parquet_mmap一节。
两种格式的核对清单
| 项目 | 流式(new_stream/open_stream) | 文件(new_file/open_file) |
|---|---|---|
| batch 数量 | 任意长度的序列 | 固定数量,读端可通过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),仅供参考