ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Arrow Python API 详解:pyarrow.fs 文件系统抽象层全解析

Apache Arrow Python API 详解:pyarrow.fs 文件系统抽象层全解析 Apache Arrow Python API 详解pyarrow.fs 文件系统抽象层全解析【免费下载链接】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 文件系统接口pyarrow.fs展开覆盖其核心接口FileSystem、FileInfo、FileSelector、六种内置实现Local、S3、GCS、HDFS、Azure、SubTree、Python 自定义文件系统PyFileSystem/FileSystemHandler/FSSpecHandler以及copy_files等实用工具。读完后你能够以统一 API 访问本地与远程存储、用 URI 自动推导文件系统类型并把 fsspec 生态无缝接入 Arrow 的数据读写链路。API 总览与模块组织pyarrow.fs是 Arrow Python 库中用于与各类本地和远程文件系统交互的抽象层。官方文档页 Filesystems 将该 API 划分为四个部分Interface接口FileInfo、FileSelector、FileSystemFilesystem Implementations文件系统实现LocalFileSystem、S3FileSystem、GcsFileSystem、HadoopFileSystem、SubTreeFileSystem、AzureFileSystemPython 侧定义的文件系统PyFileSystem、FileSystemHandler、FSSpecHandlerUtilities工具函数copy_files、initialize_s3、finalize_s3、resolve_s3_region、S3LogLevel。从源码结构看该模块的 Python 层分布在多个文件中python/pyarrow/fs.py模块入口负责按需导入各可选后端并暴露统一的公开名称python/pyarrow/_fs.pyx核心 Cython 绑定定义FileInfo、FileSelector、FileSystem、LocalFileSystem、SubTreeFileSystem、PyFileSystem、FileSystemHandler等python/pyarrow/_s3fs.pyx、python/pyarrow/_gcsfs.pyx、python/pyarrow/_azurefs.pyx、python/pyarrow/_hdfs.pyx各自云/HDFS 后端的独立绑定。这些后端的最终实现均落在 C 层的 cpp/src/arrow/filesystem/ 目录中包括 s3fs.cc、gcsfs.cc、hdfs.cc、azurefs.cc、localfs.cc 以及测试工具 mockfs.cc。值得注意的是 fs.py 中的导入策略AzureFileSystem、HadoopFileSystem、GcsFileSystem、S3FileSystem均以try/except ImportError方式可选导入失败的名字会记入_not_imported列表。当用户访问未编译进当前构建的后端时模块级__getattr__会抛出明确的错误信息# python/pyarrow/fs.py _not_imported [] try: from pyarrow._azurefs import AzureFileSystem # noqa except ImportError: _not_imported.append(AzureFileSystem) # ... def __getattr__(name): if name in _not_imported: raise ImportError( The pyarrow installation is not built with support for f{name} )也就是说若你的pyarrow安装包没有启用 S3/Azure/GCS/HDFS 支持引用对应类会得到带指引的ImportError而不是含糊的AttributeError——这是排查某个文件系统类不存在问题的第一现场。接口一FileInfo 与 FileSelectorFileInfo文件系统条目的元数据FileInfo描述一个文件系统条目文件、目录或不存在构造函数签名为FileInfo(path, typeFileType.Unknown, *, mtimeNone, mtime_nsNone, sizeNone)关键属性与取值见 _fs.pyx 中FileInfo.type的文档字符串type枚举值可为FileType.NotFound目标不存在、FileType.Unknown存在但类型未知如 Unix socket、字符设备或 Windows 的 NUL/CON 等特殊文件、FileType.File普通文件、FileType.Directory普通目录is_filetype FileType.File的便捷布尔值path文件系统内的完整路径base_name/extension路径基名与扩展名size字节大小仅对FileType.File有意义mtime/mtime_ns修改时间分别以datetime和 Unix 纪元纳秒整数表示二者互斥。一个典型的用法摘自源码文档字符串示例from pyarrow import fs local fs.LocalFileSystem() path /tmp/pyarrow-fs-example.dat with local.open_output_stream(path) as stream: stream.write(bdata) file_info local.get_file_info(path) print(file_info) # FileInfo for /tmp/pyarrow-fs-example.dat: typeFileType.File, size4 print(file_info.type) # FileType.File: 2 print(file_info.is_file) # True print(file_info.base_name) # pyarrow-fs-example.dat print(file_info.extension) # datFileSelector目录内容选择器FileSelector用于让get_file_info按条件列出目录内容构造参数为FileSelector(base_dir, allow_not_foundFalse, recursiveFalse)base_dir基准目录str 或 bytes属性可读写allow_not_found基准目录不存在时返回空列表而非抛错recursive是否递归进入子目录。示例# 递归列出目录及子目录内容 selector_1 fs.FileSelector(/tmp/alphabet, recursiveTrue) local.get_file_info(selector_1) # 只列基准目录一层 selector_2 fs.FileSelector(/tmp/alphabet) # 目录不存在时返回 [] 而不是抛异常 selector_not_found fs.FileSelector(/tmp/alphabet/missing, recursiveTrue, allow_not_foundTrue)注意get_file_info的语义约定selector 的基准目录本身不会出现在结果里不存在的普通路径不会抛错而是返回type FileType.NotFound的FileInfo——查询不到与真正的 I/O 异常被刻意区分开这在批量元数据扫描时非常实用。接口二FileSystem 基类FileSystem是抽象基类直接实例化会抛TypeError提示改用LocalFileSystem或SubTreeFileSystem等子类。其核心方法在 _fs.pyx 中定义方法说明get_file_info(paths_or_selector)接受路径、路径列表或FileSelector返回单个FileInfo或列表符号链接会被自动递归解引用create_dir(path, *, recursiveTrue)创建目录可递归建父级目录已存在时成功返回delete_dir(path)/delete_dir_contents(path, *, accept_root_dirFalse, missing_dir_okFalse)删除目录递归/ 仅删除目录内容move(src, dest)移动或重命名目标是已存在的非空目录时报错copy_file(src, dest)复制单个文件目标为目录时报错否则替换delete_file(path)删除文件open_input_file(path)打开随机访问读返回NativeFileopen_input_stream(path, compressiondetect, buffer_sizeNone)顺序读流open_output_stream(path, compressiondetect, buffer_sizeNone, metadataNone)顺序写流已存在则截断open_append_stream(path, ...)追加写部分后端不支持高效追加时会抛NotImplementedErrornormalize_path(path)规范化路径type_name/equals(other)/__eq__类型名如local、s3与相等性比较流式方法的两个细节值得注意透明压缩compression默认detect会按文件扩展名选择解压算法传具体算法名如gzip强制指定传None关闭压缩。从源码实现看_wrap_input_stream/_wrap_output_stream会在流外层按需叠加BufferedInputStream/CompressedInputStream或对应的 Output 版本即buffer_size与compression是正交的两层装饰元数据写入open_output_stream/open_append_stream支持metadata字典如Content-Type是否真正生效取决于具体后端不支持的键会被忽略。从 URI 自动构造文件系统FileSystem.from_uri(uri)是该抽象层最重要的入口之一 fs.FileSystem.from_uri(s3://usgs-landsat/collection02/) (pyarrow._s3fs.S3FileSystem object at ..., usgs-landsat/collection02)返回(FileSystem, path)元组其中path是该文件系统内部使用的抽象路径。从 _fs.pyx 的from_uri实现看原生支持的 scheme 为file、mock、s3fs、gs/gcs、hdfs/viewfs以及pathlib.Path或本地绝对路径字符串走 C 侧CFileSystemFromUriOrPath解析以fsspec前缀或hf://开头的 URI 会委托给 fsspecfsspec.url_to_fs构造PyFileSystem(FSSpecHandler(fs))——这是后文 Python 自定义文件系统的入口。C 到 Python 的类型反向映射在FileSystem.wrap中完成C 返回的type_name为local、mock、subtree、s3、gcs、abfs、hdfs、py::前缀时分别包装成对应的 Python 类。这意味着 C 层如 dataset 扫描器构造的文件系统对象在 Python 侧总能拿到类型正确的 Python 包装实例。内置文件系统实现LocalFileSystem 与 SubTreeFileSystemLocalFileSystem访问本机文件系统构造参数为LocalFileSystem(use_mmapFalse)——开启use_mmap后读文件优先走内存映射。SubTreeFileSystem(base_path, base_fs)把一个已有文件系统裁剪到某个子路径下所有操作都相对该前缀进行。两者常被组合例如把s3://bucket/dir解析出的 S3 文件系统再包一层 SubTree后续代码只需写相对路径。另外fs.py 中的内部函数_ensure_filesystem展示了这个组合的典型用法当用户以 URI 字符串形式提供文件系统且 URI 中带 path 前缀时它会先FileSystem.from_uri再校验该前缀必须是目录否则抛ValueError最后包上SubTreeFileSystem。S3FileSystem构造签名见 _s3fs.pyx 第 285 行起的__init__以关键字参数给出包含access_key、secret_key、session_token等标准 AWS 凭证字段其余参数区域、端点、重试策略等按需传入。模块级还导出了与 S3 生命周期相关的工具initialize_s3()/finalize_s3()显式初始化/终结底层 S3 会话ensure_s3_initialized()/ensure_s3_finalized()幂等版本resolve_s3_region()解析 S3 区域S3LogLevel控制 S3 相关日志级别S3RetryStrategy及其子类AwsDefaultS3RetryStrategy、AwsStandardS3RetryStrategy重试策略抽象默认策略内部维护max_attempts等参数。一个重要的实现细节在 fs.py 中S3 客户端不在导入时急切初始化而是通过ensure_s3_initialized惰性初始化并在模块导入成功时注册atexit.register(ensure_s3_finalized)。注释说明这是 GH-38364 的修复——即使程序根本没用 S3也可能在退出阶段发生崩溃因此采用用到才初始化、退出时确保终结的策略。GcsFileSystem、AzureFileSystem、HadoopFileSystem三者与 S3 类似都是可选导入 C 后端结构GcsFileSystem定义于 _gcsfs.pyx对应 C gcsfs.ccAzureFileSystem定义于 _azurefs.pyx构造以account_name为位置参数关键字参数包括account_key、blob_storage_authority等见该文件第 107 行起的__init__对应 C azurefs.ccHadoopFileSystem定义于 _hdfs.pyx对应 C hdfs.cc。C 侧每个后端都配有对应的*_test.cc如 s3fs_test.cc、azurefs_test.cc、hdfs_test.cc需要真实环境时通常使用*_test_util.cc中的测试桩或 mock 文件系统。Python 自定义文件系统PyFileSystem / FileSystemHandler / FSSpecHandler当需要行为由 Python 代码实现的文件系统时Arrow 提供三层设施FileSystemHandler抽象基类定义在 _fs.pyx 第 1333 行要求实现者提供get_type_name、get_file_info、create_dir、open_input_stream等钩子PyFileSystem把任意FileSystemHandler包装成完整的FileSystem自动获得type_namepy::handler 类型名、拷贝、移动等通用行为FSSpecHandler现成的 fsspec 适配器直接桥接fsspec.AbstractFileSystem。FssSpecHandler 的完整实现展示了桥接的具体做法例如get_file_info调用self.fs.info(path)把 fsspec 的FileNotFoundError翻译成FileInfo(path, FileType.NotFound)get_file_info_selector用self.fs.find(base_dir, maxdepth..., withdirsTrue, detailTrue)实现递归/非递归目录枚举并特意排除基准目录本身对应 GH-37555 的修复流式读写统一用pyarrow.PythonFile包装 fsspec 返回的 Python 文件对象读rb、写wb、追加ab。实际使用上fs.py 的_ensure_filesystem表明任何接受文件系统参数的 Arrow API如pq.read_table、ds.dataset等内部走该函数的接口都可以直接传一个 fsspec 实例它会被自动包成PyFileSystem(FSSpecHandler(fs))若 fsspec 实例恰好是LocalFileSystem则直接换用原生 Arrow 的LocalFileSystem。此外也可显式走 URI 形式 fs.FileSystem.from_uri(fsspecmemory:///path/to/file) (pyarrow._fs.PyFileSystem object at ..., /path/to/file)工具函数copy_files 与 S3 辅助函数copy_files跨文件系统复制copy_files(source, destination, source_filesystemNone, destination_filesystemNone, *, chunk_size1024*1024, use_threadsTrue)支持在不同文件系统之间递归复制文件与目录例如从 S3 拉到本地。参数语义源自 fs.py 的文档字符串source/destination路径或 URI。source为目录时递归复制source为文件时destination解释为目标文件名而非目录需要时自动创建目录source_filesystem/destination_filesystem非 URI 情形下必须显式指定URI 情形自动推导chunk_size默认 1MB单次读块上限加大可适配高延迟远程文件系统但更耗内存use_threads默认True多线程加速复制。实现上copy_files先经_resolve_filesystem_and_path解析两端该函数会优先尝试把路径当作本地相对路径本地不存在时再按 URI 解析既不是有效 URI 也不存在本地文件时抛出更友好的 file-not-found 错误然后判断source是目录还是文件分别落到 C 扩展层的_copy_files_selector/_copy_files定义于 _fs.pyx 第 1617 行、1642 行附近执行分块拷贝。官方文档字符串给出的示例 s3, path fs.FileSystem.from_uri( ... s3://registry.opendata.aws/roda/ndjson/) selector fs.FileSelector(path) s3.get_file_info(selector) [FileInfo for registry.opendata.aws/roda/ndjson/index.ndjson:...] fs.copy_files(s3://registry.opendata.aws/roda/ndjson/index.ndjson, ... ffile:///{local_path}/index_copy.ndjson) # 用显式 FileSystem 对象的形式 fs.copy_files(registry.opendata.aws/roda/ndjson/index.ndjson, ... ffile:///{local_path}/index_copy.ndjson, ... source_filesystemfs.S3FileSystem())S3 辅助工具小结文档 Utilities 一节列出的initialize_s3、finalize_s3、resolve_s3_region、S3LogLevel与上述S3RetryStrategy家族共同构成 S3 后端的运维面显式控制底层会话生命周期推荐依赖惰性初始化 atexit 自动终结的默认行为即可、按 endpoint 解析区域、调整日志级别、以及定制重试行为。由于这些符号都来自可选模块 python/pyarrow/_s3fs.pyx未启用 S3 支持的构建中引用它们会触发前文所述的ImportError。小结如何把 pyarrow.fs 用进数据管线本地/简单场景直接使用LocalFileSystem()或干脆传字符串路径让 Arrow 内部_resolve_filesystem_and_path自动推导本地存在→本地文件系统否则按 URI 解析云存储场景优先fs.FileSystem.from_uri(s3://...)/gs:///az://等 URI 形式让 scheme 决定后端注意pyarrow安装时是否启用了相应后端缺失时ImportError会明确点名缺哪个类前缀裁剪用SubTreeFileSystem把一个 bucket 内的目录暴露为根简化下游相对路径生态兼容fsspec 用户可以直接把自己的 fsspec 实例传给 Arrow 各 API或写fsspec/hf://前缀 URI有完全自定义需求时继承FileSystemHandler并用PyFileSystem包装即可。文档与实现的关键参照docs/source/python/api/filesystems.rst、python/pyarrow/fs.py、python/pyarrow/_fs.pyx、python/pyarrow/_s3fs.pyx、python/pyarrow/_gcsfs.pyx、python/pyarrow/_azurefs.pyx、python/pyarrow/_hdfs.pyx 以及 C 后端目录 cpp/src/arrow/filesystem/。【免费下载链接】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),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进