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_image、guess_mime_type、download等表达式做后续的 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()S3Config与GCSConfig均由 daft/io/init.py 导出,read_blob同样在 daft/io/init.py 中作为公共 API 公开。若未显式传入io_config,read_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:各参数说明:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
path | str \| list[str] | 必填 | 文件路径或路径列表,支持s3://、gs://等远程 URI 及通配符;传入空列表会抛出ValueError(源码中显式校验,见 daft/io/_blob.py) |
max_connections | int | 32 | 每个线程用于下载文件内容的连接数上限,用于控制并发下载的吞吐量 |
on_error | "raise" \| "null" | "raise" | 下载出错时的行为:"raise"立即抛错,"null"记录错误并回退为null的content值 |
io_config | IOConfig \| None | None | 访问远程对象存储时的配置;为None时使用 Daft 计划配置中的默认 IO 配置 |
关于max_connections的调优,daft/functions/url.py 中download的文档给出了实用建议:如果遇到过多 S3 错误(超时、DNS 错误、减速错误),可以调低max_connections以减轻对 S3 服务端的压力;如果机器核数较少但网络带宽很高,则可以调高它以获得更高的吞吐量。
输出 Schema
read_blob返回的 DataFrame 固定包含三列:
| 列名 | 类型 | 描述 |
|---|---|---|
path | string | 文件路径 |
size | int64 | 文件大小(字节) |
content | binary | 文件的原始字节内容 |
该 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())]),并验证单文件场景下行数、路径后缀、大小与内容字节均正确。空文件同样能正常读入(size为0、content为空字节串),见同一文件的test_read_blob_empty_file。
Glob 通配规则
read_blob支持与 Daft 其他读取器一致的通配符语法,这一点在 daft/io/_blob.py 与 daft/io/file_path.py 的文档字符串中均有明确记载:
*匹配任意数量(含零个)的任意字符?匹配任意单个字符[...]匹配括号中的任意单个字符**递归匹配任意层级的目录
# 目录下所有 .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 模式列出文件,产出包含path、size以及(针对 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_path与download的组合,需要精细化控制(如先过滤再下载)时可直接使用底层 API; - 配合
decode_image、guess_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),仅供参考