Daft read_blob 使用指南:将任意文件作为二进制 Blob 读入 DataFrame
2026/9/17 22:49:57 网站建设 项目流程

Daft read_blob 使用指南:将任意文件作为二进制 Blob 读入 DataFrame

【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft

导读

daft.read_blob()是 Daft 提供的高层便捷 API,用于把磁盘或对象存储(S3、GCS 等)上的任意文件以原始字节形式读入 DataFrame,每个文件对应一行。它非常适合处理非表格类数据——图片、音频、PDF、HDF5 或任何二进制文件——让这些数据直接进入 Daft 的列式执行引擎,配合decode_imageguess_mime_typedownload等表达式做后续的 AI/多模态处理。读完本文,你将掌握read_blob的完整参数、输出 Schema、glob 通配规则、错误处理策略,以及它与from_glob_path+download的关系,并能在本地与 S3/GCS 场景下直接复现运行。

基本用法:本地与远程对象存储

read_blob接受一个路径字符串或路径列表(支持通配符),每个匹配文件被读为一行。它的定位与 DuckDB 的read_blob类似,这在 daft/io/_blob.py 的函数文档字符串中有明确说明。

本地文件

import daft df = daft.read_blob("/path/to/files/*.jpeg") df.show()

S3 上的远程文件

import daft from daft.io import S3Config, IOConfig io_config = IOConfig(s3=S3Config(region_name="us-west-2", anonymous=True)) df = daft.read_blob("s3://my-bucket/images/*.jpeg", io_config=io_config) df.show()

GCS 上的远程文件

import daft from daft.io import GCSConfig, IOConfig io_config = IOConfig(gcs=GCSConfig(anonymous=True)) df = daft.read_blob("gs://my-bucket/images/*.jpeg", io_config=io_config) df.show()

S3ConfigGCSConfig均由 daft/io/init.py 导出,read_blob同样在 daft/io/init.py 中作为公共 API 公开。若未显式传入io_configread_blob会回退到get_context().daft_planning_config.default_io_config(见 daft/io/_blob.py),即沿用 Daft 的默认 IO 配置。

完整参数签名

从 daft/io/_blob.py 的源码可以看到read_blob的完整签名:

def read_blob( path: str | list[str], *, max_connections: int = 32, on_error: Literal["raise", "null"] = "raise", io_config: IOConfig | None = None, ) -> DataFrame:

各参数说明:

参数类型默认值说明
pathstr \| list[str]必填文件路径或路径列表,支持s3://gs://等远程 URI 及通配符;传入空列表会抛出ValueError(源码中显式校验,见 daft/io/_blob.py)
max_connectionsint32每个线程用于下载文件内容的连接数上限,用于控制并发下载的吞吐量
on_error"raise" \| "null""raise"下载出错时的行为:"raise"立即抛错,"null"记录错误并回退为nullcontent
io_configIOConfig \| NoneNone访问远程对象存储时的配置;为None时使用 Daft 计划配置中的默认 IO 配置

关于max_connections的调优,daft/functions/url.py 中download的文档给出了实用建议:如果遇到过多 S3 错误(超时、DNS 错误、减速错误),可以调低max_connections以减轻对 S3 服务端的压力;如果机器核数较少但网络带宽很高,则可以调高它以获得更高的吞吐量。

输出 Schema

read_blob返回的 DataFrame 固定包含三列:

列名类型描述
pathstring文件路径
sizeint64文件大小(字节)
contentbinary文件的原始字节内容

该 Schema 由 tests/io/test_read_blob.py 中的测试test_read_blob_single_file直接断言验证:它通过Schema.from_pyarrow_schema对比了pa.schema([("path", pa.string()), ("size", pa.int64()), ("content", pa.binary())]),并验证单文件场景下行数、路径后缀、大小与内容字节均正确。空文件同样能正常读入(size0content为空字节串),见同一文件的test_read_blob_empty_file

Glob 通配规则

read_blob支持与 Daft 其他读取器一致的通配符语法,这一点在 daft/io/_blob.py 与 daft/io/file_path.py 的文档字符串中均有明确记载:

  1. *匹配任意数量(含零个)的任意字符
  2. ?匹配任意单个字符
  3. [...]匹配括号中的任意单个字符
  4. **递归匹配任意层级的目录
# 目录下所有 .png 文件 df = daft.read_blob("/data/*.png") # 递归搜索所有 .png 文件 df = daft.read_blob("/data/**/*.png") # 显式指定多个路径 df = daft.read_blob(["/data/a.bin", "/data/b.bin"])

通配符匹配的测试行为可参考 tests/io/test_read_blob.py 的test_read_blob_glob(3 个.bin文件逐行校验内容与大小)以及test_read_blob_multiple_paths(验证显式多路径列表输入)。

错误处理:raise 与 null 两种策略

默认情况下,任何文件的下载错误都会立即抛出异常(on_error="raise")。如果希望容错地跳过坏文件,可以设置on_error="null",此时错误会被记录(log),而该行的content回退为null

df = daft.read_blob("/data/*.bin", on_error="null")

这一行为同样适用于底层download表达式:on_error的语义在 daft/functions/url.py 中有完整描述——"raise"立即抛出错误,"null"记录错误但回退为null值,返回的表达式是二进制类型,出错时为None

常见用例

解码图片

读取图片文件并解码为可做下游处理的图像列:

import daft from daft.functions import decode_image df = daft.read_blob("s3://my-bucket/images/*.jpeg", io_config=io_config) df = df.with_column("image", decode_image(df["content"]))

decode_image位于 daft/functions/image.py,只接受包含已编码图像(如 PNG、JPEG)的二进制列,并支持两个可调参数:on_error(默认"raise",可设为"null"以容错)和mode(默认ImageMode.RGB,即统一转为 RGB 存储;设为None时从原始数据推断模式)。解码后还可以继续接image_to_tensor(daft/functions/image.py)将图像转为张量,方便送入深度学习模型。

检测文件类型

通过读取文件的 magic bytes(魔数)判断每个文件的 MIME 类型:

import daft from daft.functions import guess_mime_type df = daft.read_blob("/data/**/*") df = df.with_column("mime", guess_mime_type(df["content"]))

guess_mime_type(daft/functions/file_.py)通过检查二进制数据的 magic bytes 来识别常见格式,包括 PNG、JPEG、GIF、WEBP、PDF、ZIP、MP3、WAV、OGG、MP4、MPEG、HDF5 和 HTML;无法识别时返回None。其文档中还给出了一个可直接运行的示例:输入b"\x89PNG\r\n\x1a\n"会得到"image/png",输入未知字节则得到None

与其他 API 的关系:从底层组合出 read_blob

read_blob本身是一个组合式的便捷 API。从 daft/io/_blob.py 的源码可见其完整实现:

df = from_glob_path(path, io_config=io_config) return df.select( col("path"), col("size"), download( col("path"), max_connections=max_connections, on_error=on_error, io_config=io_config, ).alias("content"), )

即它由两个底层构件拼接而成:

  • daft.from_glob_path()(daft/io/file_path.py):负责按 glob 模式列出文件,产出包含pathsize以及(针对 parquet 的)rows列的元数据 DataFrame。底层通过LogicalPlanBuilder.from_glob_scan构建逻辑计划。注意:如果 glob 模式没有匹配到任何文件,返回的是空 DataFrame 而不是报错(见 daft/io/file_path.py 的 Note)。
  • daft.functions.download(daft/functions/url.py):把path字符串列当作 URL,逐行下载字节内容,生成一个二进制表达式。

因此,当需要更细粒度的控制——例如只下载过滤后的文件子集——可以直接使用这两个构件。典型的组合方式如下:

import daft from daft.functions import download # 1. 先只列出文件元数据(不下载内容) df = daft.from_glob_path("s3://my-bucket/images/*.jpeg", io_config=io_config) # 2. 按元数据过滤,例如只保留大于 1KB 的文件 df = df.where(df["size"] > 1024) # 3. 再只对过滤后的子集执行下载 df = df.with_column("content", download(df["path"], io_config=io_config))

从实现细节看,download在底层会把max_connections覆盖到io_config的 S3 配置中(Rust 层实际执行min(S3Config.max_connections, url_download.max_connections),见 daft/functions/url.py),并依据执行模式选择多线程或单线程的 tokio IO 运行时:Ray 分布式执行时默认单线程运行时以避免海量连接,本地执行时则共享多线程运行时的连接池(daft/functions/url.py)。

小结

  • daft.read_blob()以"一行一文件"的方式把任意二进制文件读入 DataFrame,输出固定为path/size/content三列;
  • 支持本地路径、s3://gs://*?[...]**四种通配符,也支持传入路径列表;
  • 通过on_error="null"可对下载失败的单个文件做容错处理;
  • 它本质上是from_glob_pathdownload的组合,需要精细化控制(如先过滤再下载)时可直接使用底层 API;
  • 配合decode_imageguess_mime_type等表达式,可以快速搭建图片解码、文件类型识别等 AI/多模态数据处理流水线。

相关参考:实现源码 daft/io/_blob.py、daft/io/file_path.py、daft/functions/url.py;功能函数 daft/functions/image.py、daft/functions/file_.py;单元测试 tests/io/test_read_blob.py;S3/MinIO 集成测试 tests/integration/io/test_blob_s3_minio.py。

【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询