
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与文件访问的核心 APIinput_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 行附近各派生类分布在第 8811848 行工厂函数则定义在第 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 TrueBufferOutputStream仅置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, compressiondetect, buffer_sizeNone)参数说明source可以是 str、Path、buffer 或 file-like 对象compression默认detect若 source 是文件路径则按文件扩展名推断压缩算法传None表示不应用压缩否则必须指定受支持的算法名如gzipbuffer_sizeNone或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(bsome data) with pa.input_stream(buf) as stream: stream.read(4) # - bsome # 2) 从 gzip 文件路径创建 OSFile 自动解压 with gzip.open(example.gz, wb) as f: f.write(bsome data) with pa.input_stream(example.gz) as stream: stream.read() # - bsome data # 3) 从普通文本文件路径 with open(example.txt, modew) as f: f.write(some text) with pa.input_stream(example.txt) as stream: stream.read(6) # - bsome t3.2 output_stream对称的输出生成器pa.output_stream(source, compressiondetect, buffer_sizeNone)参数语义与input_stream完全对称见 io.pxi 第 2869 行起。分派差异在于对 buffer 来源输出方向选用的是FixedSizeBufferWriter而非BufferReaderelif isinstance(source, (Buffer, memoryview)): stream FixedSizeBufferWriter(as_buffer(source))这与 BufferOutputStream 形成对比后者面向可增长的 resizable buffer写完后用getvalue()取回结果而FixedSizeBufferWriter面向预先分配好大小的 buffer。文档示例import pyarrow as pa data bbuffer 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) # - bbuffer # 文件路径则直接得到 OSFile(w) with pa.output_stream(example_second.txt) as stream: stream.write(bWrite some data) # - 153.3 memory_map 与 create_memory_map内存映射memory_mapdef memory_map(path, moder): Open memory map at file path. Size of the memory map cannot change.mode取值为r读、r读写、w写默认r源码中非法模式会抛出ValueError(fInvalid 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(bConstructing 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(bCreate a memory-mapped file) mmap.read_at(10, 9) # - bmemory-map对应的MemoryMappedFile类io.pxi 第 1013 行起额外提供resize(new_size)方法——同时调整映射与底层文件大小——以及fileno()返回底层文件描述符。四、流类逐个解析源码实现要点4.1 OSFile普通文件描述符后端OSFileio.pxi 第 1184 行是“backed by a regular file descriptor”的流构造参数path既接受路径字符串也接受已打开的文件描述符传入的 fd 归 OSFile 所有并由其关闭。文档示例同时展示了能力探测with pa.OSFile(example_osfile.arrow, modew) as f: f.writable() # - True f.write(bOSFile) # - 6 f.seekable() # - False值得注意的是seekable()返回False以w模式新建的文件流并非随机访问文件这与MemoryMappedFile支持read_at形成鲜明对照选流时应根据是否要随机定位来决定。4.2 PythonFile包装 Python file-like 对象PythonFileio.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 breader data buf memoryview(data) with pa.input_stream(buf) as stream: stream.size() # - 11 stream.read(6) # - breader stream.seek(7) stream.read(15) # - bdata4.4 BufferOutputStream 与 FixedSizeBufferWriter两种内存写端二者都写内存但生命周期模型不同BufferOutputStreamio.pxi 第 1695 行写入一个可增长resizablebuffer构造时可传MemoryPool数据不落地直到调用getvalue()——该方法内部先Close()流nogil下执行再把 buffer 包装返回并流进入已关闭状态f pa.BufferOutputStream() f.write(bpyarrow.Buffer) # - 14 f.getvalue() # pyarrow.Buffer address... size14 is_cpuTrue is_mutableTrue它是“先写后取”模式的事实标准文档中CompressedOutputStream的示例也以其作为底层承载。FixedSizeBufferWriterio.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(babcde) # - 5 buf.to_pybytes() # babcde4.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 bCompressed stream raw pa.BufferOutputStream() with pa.CompressedOutputStream(raw, gzip) as compressed: compressed.write(data) # - 17 cdata raw.getvalue() # 方式一工厂函数自动识别并包裹 with pa.input_stream(cdata, compressiongzip) as compressed: compressed.read() # - bCompressed stream # 方式二手工组合 BufferReader CompressedInputStream raw pa.BufferReader(cdata) with pa.CompressedInputStream(raw, gzip) as compressed: compressed.read() # - bCompressed stream该示例恰好印证了input_stream(source, compression...)的内部行为——它等价于BufferReaderCompressedInputStream的组合即第三节中“工厂函数按需包裹装饰器”的具体体现。4.6 各流类能力速查流类方向底层承载随机访问源码位置NativeFile基类C 流句柄可选io.pxi#L101OSFile读/写/读写文件描述符视模式io.pxi#L1184PythonFile读/写Python file-like视对象io.pxi#L881BufferReader读内存 buffer零拷贝支持io.pxi#L1751BufferOutputStream写可增长 buffergetvalue()取回不支持io.pxi#L1695FixedSizeBufferWriter写固定大小 buffer不支持io.pxi#L1315MemoryMappedFile读/写内存映射文件支持read_at/resizeio.pxi#L1013CompressedInputStream读任意输入流 Codec继承内层io.pxi#L1791CompressedOutputStream写任意输出流 Codec—io.pxi#L1848五、实践选型工厂函数优先流类兜底结合文档定位与源码实现可以给出如下决策路径日常读写一律先用工厂函数。pa.input_stream/pa.output_stream屏蔽了来源类型判断并统一提供compressiondetect按扩展名识别 gzip 等格式与buffer_size缓冲选项。源码中“source 已是NativeFile时直接透传”的分支io.pxi#L2842-L2845也说明工厂函数与手写流类可以无缝混用。需要随机访问或零拷贝时用pa.memory_map/pa.create_memory_map或直接构造BufferReaderMemoryMappedFile的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可读/可写/可选 seekwith语义关闭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),仅供参考