PyArrow 流与文件访问 API 详解:从 input_stream 工厂函数到 NativeFile 流类体系
【免费下载链接】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
本文基于 Apache Arrow 官方 Python API 文档 Streams and File Access 参考页,系统讲解 PyArrow 中流(Stream)与文件访问的核心 API:input_stream/output_stream/memory_map/create_memory_map四个推荐工厂函数,以及NativeFile、OSFile、PythonFile、BufferReader、BufferOutputStream、FixedSizeBufferWriter、MemoryMappedFile、CompressedInputStream、CompressedOutputStream等流类。结合 python/pyarrow/io.pxi 中的 Cython 源码实现,你既能掌握每个 API 的参数语义与可复制的用法示例,也能理解其底层的流构建链路与压缩、内存映射机制,最终形成从“选择正确工厂函数”到“定制底层流类”的完整实战能力。
一、API 总览:工厂函数、流类与文件系统三层结构
官方参考页 docs/source/python/api/files.rst 将“流与文件访问”划分为三个层次:
- Factory Functions(工厂函数):
input_stream、output_stream、memory_map、create_memory_map。文档明确指出,这些工厂函数是创建 Arrow 流的推荐方式,它们接受多种来源(in-memory buffers 或 on-disk files),并自动选择合适的流类实现。 - Stream Classes(流类):
NativeFile(基类)、OSFile、PythonFile、BufferReader、BufferOutputStream、FixedSizeBufferWriter、MemoryMappedFile、CompressedInputStream、CompressedOutputStream。 - File Systems:参考页将文件系统部分指向 Filesystem Interface 文档(
:ref: api.fs),即pyarrow.fs模块,提供LocalFileSystem、S3FileSystem、GcsFileSystem等实现。
从源码结构看,上述全部内容都集中在 python/pyarrow/io.pxi(约 2957 行)中实现:工厂函数与流类定义均在该文件内,NativeFile基类位于第 101 行附近,各派生类分布在第 881~1848 行,工厂函数则定义在第 1113、1154、2783、2869 行附近。
二、NativeFile 基类:所有 Arrow 流的统一抽象
所有流类都继承自NativeFile。根据 io.pxi 中的类注释,其语义可归纳为三点:
- 能力声明:流要么可读(readable)、要么可写(writable)、要么两者兼有,并且可选地支持 seek。基类在
__cinit__中初始化own_file、is_readable、is_writable、is_seekable、_is_appending等标志位,各子类在构造时按需打开对应标志(如BufferReader仅置is_readable = True,BufferOutputStream仅置is_writable = True)。 - 主要用途:类注释强调,虽然
NativeFile暴露了 Python 层的读写方法,但其首要意图是被传递给其他 Arrow 组件(如 Arrow IPC 读写例程)使用,而不是当作普通文件反复手动读写。 - 与普通 Python 文件的关键差异:注释中特别警示——销毁一个可写的 Arrow 流而不显式关闭,不会 flush 任何 pending 数据。这是与普通 Python 文件(GC 时自动 flush)最重要的行为差异,生产代码应始终使用
with语句或显式close()。
基类还实现了上下文管理器协议(__enter__/__exit__,退出时自动close())和mode属性,后者根据可读/可写标志位模拟内建文件模式:rb、wb、rb+、ab。另外,基类定义了_default_chunk_size = 256 * 1024,注释说明该分块大小是为网络文件系统特意取大值,供分块读取使用。
三、工厂函数:创建流的推荐入口
3.1 input_stream:多来源自适应的输入流
签名与参数(源码位置:io.pxi 第 2783 行):
pa.input_stream(source, compression='detect', buffer_size=None)| 参数 | 说明 |
|---|---|
source | 可以是 str、Path、buffer 或 file-like 对象 |
compression | 默认'detect':若 source 是文件路径,则按文件扩展名推断压缩算法;传None表示不应用压缩;否则必须指定受支持的算法名(如"gzip") |
buffer_size | None或0表示不缓冲;否则作为临时读缓冲区的字节数 |
源码中的分派逻辑非常清晰,展示了工厂函数“自动选类”的实现:
if isinstance(source, NativeFile): stream = source # 已是 Arrow 流,直接透传 elif source_path is not None: stream = OSFile(source_path, 'r') # 文件路径 → OSFile elif isinstance(source, (Buffer, memoryview)): stream = BufferReader(as_buffer(source)) # 内存缓冲 → BufferReader elif (hasattr(source, 'read') and hasattr(source, 'close') and hasattr(source, 'closed')): stream = PythonFile(source, 'r') # file-like → PythonFile随后按需包裹两层装饰器:buffer_size非空时外包BufferedInputStream,压缩非None时再外包CompressedInputStream。这解释了文档示例中为什么读取example.gz时无需手动解压缩——'detect'通过扩展名解析出 gzip 后自动完成了包裹。
文档自带的可复制示例:
import pyarrow as pa import gzip # 1) 从内存 buffer 创建 BufferReader buf = memoryview(b"some data") with pa.input_stream(buf) as stream: stream.read(4) # -> b'some' # 2) 从 gzip 文件路径创建 OSFile + 自动解压 with gzip.open('example.gz', 'wb') as f: f.write(b'some data') with pa.input_stream('example.gz') as stream: stream.read() # -> b'some data' # 3) 从普通文本文件路径 with open('example.txt', mode='w') as f: f.write('some text') with pa.input_stream('example.txt') as stream: stream.read(6) # -> b'some t'3.2 output_stream:对称的输出生成器
pa.output_stream(source, compression='detect', buffer_size=None)参数语义与input_stream完全对称(见 io.pxi 第 2869 行起)。分派差异在于:对 buffer 来源,输出方向选用的是FixedSizeBufferWriter而非BufferReader:
elif isinstance(source, (Buffer, memoryview)): stream = FixedSizeBufferWriter(as_buffer(source))这与 BufferOutputStream 形成对比:后者面向可增长的 resizable buffer,写完后用getvalue()取回结果;而FixedSizeBufferWriter面向预先分配好大小的 buffer。文档示例:
import pyarrow as pa data = b"buffer data" empty_obj = bytearray(11) buf = pa.py_buffer(empty_obj) with pa.output_stream(buf) as stream: stream.write(data) # -> 11 with pa.input_stream(buf) as stream: stream.read(6) # -> b'buffer' # 文件路径则直接得到 OSFile('w') with pa.output_stream('example_second.txt') as stream: stream.write(b'Write some data') # -> 153.3 memory_map 与 create_memory_map:内存映射
memory_map:
def memory_map(path, mode='r'): """Open memory map at file path. Size of the memory map cannot change."""mode取值为'r'(读)、'r+'(读写)、'w'(写),默认'r';源码中非法模式会抛出ValueError(f'Invalid file mode: {mode}')。- 底层调用 C++ 侧
CMemoryMappedFile.Open,返回对象同时挂接 output stream 与 random access file 两种句柄,因此支持随机访问(read_at)。 - 关键约束:打开后内存映射的大小不可变,且
_check_is_file会拒绝目录路径。
文档示例展示了零内存拷贝的读取:
import pyarrow as pa with pa.output_stream('example_mmap.txt') as stream: stream.write(b'Constructing a buffer referencing the mapped memory') # 51 with pa.memory_map('example_mmap.txt') as mmap: mmap.read_at(6, 45) # 从偏移 6 读 45 字节create_memory_map 则用于先创建再映射:
def create_memory_map(path, size): """Create a file of the given size and memory-map it.""" return MemoryMappedFile.create(path, size)示例:
with pa.create_memory_map('example_mmap_create.dat', 27) as mmap: mmap.write(b'Create a memory-mapped file') mmap.read_at(10, 9) # -> b'memory-map'对应的MemoryMappedFile类(io.pxi 第 1013 行起)额外提供resize(new_size)方法——同时调整映射与底层文件大小——以及fileno()返回底层文件描述符。
四、流类逐个解析:源码实现要点
4.1 OSFile:普通文件描述符后端
OSFile(io.pxi 第 1184 行)是“backed by a regular file descriptor”的流,构造参数path既接受路径字符串也接受已打开的文件描述符,传入的 fd 归 OSFile 所有并由其关闭。文档示例同时展示了能力探测:
with pa.OSFile('example_osfile.arrow', mode='w') as f: f.writable() # -> True f.write(b'OSFile') # -> 6 f.seekable() # -> False值得注意的是seekable()返回False:以'w'模式新建的文件流并非随机访问文件,这与MemoryMappedFile支持read_at形成鲜明对照,选流时应根据是否要随机定位来决定。
4.2 PythonFile:包装 Python file-like 对象
PythonFile(io.pxi 第 881 行)用于桥接任意符合 file-like 协议的 Python 对象(具备read/write/close/closed)。input_stream对“非路径、非 buffer”的来源即落到此类。适合接入io.BytesIO、自定义网络句柄等,代价是经由 Python 层转发、无法利用底层 C++ 的直接系统调用路径。
4.3 BufferReader:零拷贝内存读取
BufferReader 的类注释即“Zero-copy reader from objects convertible to Arrow buffer”。构造时把bytes或pyarrow.Buffer转成Buffer后直接new CBufferReader(...)挂为 random access file,因此size()、seek()、read_at()全部可用:
data = b'reader data' buf = memoryview(data) with pa.input_stream(buf) as stream: stream.size() # -> 11 stream.read(6) # -> b'reader' stream.seek(7) stream.read(15) # -> b'data'4.4 BufferOutputStream 与 FixedSizeBufferWriter:两种内存写端
二者都写内存,但生命周期模型不同:
BufferOutputStream(io.pxi 第 1695 行):写入一个可增长(resizable)buffer,构造时可传
MemoryPool;数据不落地,直到调用getvalue()——该方法内部先Close()流(nogil下执行),再把 buffer 包装返回,并流进入已关闭状态:f = pa.BufferOutputStream() f.write(b'pyarrow.Buffer') # -> 14 f.getvalue() # <pyarrow.Buffer address=... size=14 is_cpu=True is_mutable=True>它是“先写后取”模式的事实标准,文档中
CompressedOutputStream的示例也以其作为底层承载。FixedSizeBufferWriter(io.pxi 第 1315 行):面向固定大小buffer 的写端,底层是 C++ 的
CFixedSizeBufferWriter;当目标是预分配的 Arrow buffer(如pa.allocate_buffer(n)或外部共享内存区)时选用。源码还暴露了set_memcopy_threads与set_memcopy_blocksize两个调优接口,用于配置内存拷贝的线程数与块大小,适合大批量写入场景。buf = pa.allocate_buffer(5) with pa.output_stream(buf) as stream: stream.write(b'abcde') # -> 5 buf.to_pybytes() # b'abcde'
4.5 CompressedInputStream / CompressedOutputStream:透明压缩包装
这两个类是流包装器而非独立存储后端。CompressedInputStream 的构造签名要求compression非None,支持的算法为"bz2"、"brotli"、"gzip"、"lz4"、"zstd";其实现先经get_native_file把任意来源归一化为NativeFile,再用Codec与CCompressedInputStream.Make(codec, reader)完成 C++ 侧包装,实现边读边解压。输出侧对称:CCompressedOutputStream.Make包裹任意写端,实现边写边压缩。
文档示例完整演示了“压缩写入 → 解压读回”的闭环:
import pyarrow as pa data = b"Compressed stream" raw = pa.BufferOutputStream() with pa.CompressedOutputStream(raw, "gzip") as compressed: compressed.write(data) # -> 17 cdata = raw.getvalue() # 方式一:工厂函数自动识别并包裹 with pa.input_stream(cdata, compression="gzip") as compressed: compressed.read() # -> b'Compressed stream' # 方式二:手工组合 BufferReader + CompressedInputStream raw = pa.BufferReader(cdata) with pa.CompressedInputStream(raw, "gzip") as compressed: compressed.read() # -> b'Compressed stream'该示例恰好印证了input_stream(source, compression=...)的内部行为——它等价于BufferReader+CompressedInputStream的组合,即第三节中“工厂函数按需包裹装饰器”的具体体现。
4.6 各流类能力速查
| 流类 | 方向 | 底层承载 | 随机访问 | 源码位置 |
|---|---|---|---|---|
NativeFile | 基类 | C++ 流句柄 | 可选 | io.pxi#L101 |
OSFile | 读/写/读写 | 文件描述符 | 视模式 | io.pxi#L1184 |
PythonFile | 读/写 | Python file-like | 视对象 | io.pxi#L881 |
BufferReader | 读 | 内存 buffer(零拷贝) | 支持 | io.pxi#L1751 |
BufferOutputStream | 写 | 可增长 buffer,getvalue()取回 | 不支持 | io.pxi#L1695 |
FixedSizeBufferWriter | 写 | 固定大小 buffer | 不支持 | io.pxi#L1315 |
MemoryMappedFile | 读/写 | 内存映射文件 | 支持(read_at/resize) | io.pxi#L1013 |
CompressedInputStream | 读 | 任意输入流 + Codec | 继承内层 | io.pxi#L1791 |
CompressedOutputStream | 写 | 任意输出流 + Codec | — | io.pxi#L1848 |
五、实践选型:工厂函数优先,流类兜底
结合文档定位与源码实现,可以给出如下决策路径:
- 日常读写一律先用工厂函数。
pa.input_stream/pa.output_stream屏蔽了来源类型判断,并统一提供compression='detect'(按扩展名识别 gzip 等格式)与buffer_size缓冲选项。源码中“source 已是NativeFile时直接透传”的分支(io.pxi#L2842-L2845)也说明工厂函数与手写流类可以无缝混用。 - 需要随机访问或零拷贝时,用
pa.memory_map/pa.create_memory_map或直接构造BufferReader;MemoryMappedFile的read_at(offset, length)对大文件局部读取尤其有用,且读操作不经过内存分配与拷贝。 - 需要透明压缩时,优先把
compression参数交给工厂函数;只有在需要复用自定义流链(例如先经过BufferedInputStream再压缩)时才手工组合CompressedInputStream/CompressedOutputStream。 - 可写流的关闭纪律:由于
NativeFile不 flush on GC(类注释),所有可写流都应放在with块中——这一点在BufferOutputStream.getvalue()内部会显式Close()的行为上也能得到印证。
六、与文件系统层的衔接
参考页末尾将 File Systems 指向 docs/source/python/filesystems.rst。该文档说明pyarrow.fs提供抽象FileSystem基类及LocalFileSystem、S3FileSystem、GcsFileSystem、HadoopFileSystem、AzureFileSystem等实现,并支持 fsspec 兼容文件系统。其输出的open_input_stream/open_output_stream返回的同样是NativeFile体系的对象,因此本文介绍的工厂函数、压缩包装与流类能力可以无差别地作用于任何文件系统来源的流。例如 filesystems.rst 中的示例:
import pyarrow as pa table = pa.table({'col1': [1, 2, 3]}) local = fs.LocalFileSystem() with local.open_output_stream("test.arrow") as file: with pa.RecordBatchFileWriter(file, table.schema) as writer: writer.write_table(table)这里RecordBatchFileWriter接收的正是NativeFile——印证了基类注释所说“流的首要意图是传递给其他 Arrow 组件”。
七、小结
PyArrow 的文件与流访问体系可以概括为“一个基类、四个工厂、两类能力”:
- 一切流继承自
NativeFile(可读/可写/可选 seek,with语义关闭,GC 不 flush); - 四个工厂函数(
input_stream、output_stream、memory_map、create_memory_map)负责来源归一化、压缩检测与缓冲包裹,是首选入口; - 具体流类各自覆盖文件描述符、Python 对象、内存 buffer(定长/增长)、内存映射与透明压缩五类场景,其参数与行为均能在 python/pyarrow/io.pxi 中找到一一对应的 Cython/C++ 实现。
掌握这套体系后,无论是本地文件、内存数据还是远端存储流,都可以用统一的方式接入 Parquet、IPC、Feather 等所有 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
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考