Daft 读取 WebDataset:以 TAR 分片为载体的多模态数据引擎接入指南
【免费下载链接】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
WebDataset 是一种将多模态样本按“样本前缀 + 后缀”连续打包进 TAR 分片(shard)的存储格式,广泛用于大规模图文、音视频训练数据集的组织。本文基于 Daft 源码库,完整讲解如何通过daft.read_webdataset()与daft.read_huggingface()读取这类数据,覆盖路径传参、列语义、惰性媒体读取、Hugging Face 集成、格式约束与底层实现原理,帮助读者在 Daft 中直接对 WebDataset 数据执行投影、筛选与分析,而不必先解包整个归档。
WebDataset 在 Daft 中的定位
WebDataset 将一张图片及其对应的 JSON 标注、文本说明、类别标签等“附属文件”作为同一 TAR 归档中的连续成员存储,成员文件名遵循prefix.suffix约定。Daft 将每个prefix(样本前缀)视为一行数据,把suffix(后缀)映射为列名,从而把“一个归档文件”解构为“一张多模态表格”。Daft 的核心入口在 daft/io/webdataset/_webdataset.py 的read_webdataset(),并由 daft/io/webdataset/init.py 导出。
import daft df = daft.read_webdataset("/datasets/images/*.tar")读取 TAR 分片:路径的四种形式
read_webdataset的第一个参数path可以接受四种形式(对应源码 daft/io/webdataset/_webdataset.py 中_resolve_archive_paths/_glob_pattern的逻辑):
| 形式 | 示例 | 说明 |
|---|---|---|
| 单个 TAR 文件 | "/datasets/images/shard-000.tar" | 直接读取该分片 |
| 目录 | "/datasets/images" | 递归搜索目录下所有*.tar文件(源码中会补全为**/*.tar) |
| glob 通配符 | "/datasets/images/*.tar" | 匹配所有分片 |
| 路径列表 | ["/a/1.tar", "/s3://bucket/2.tar"] | 混合本地与远程存储 |
路径解析后仅保留类型为File且以.tar结尾的结果;若一个分片被多个模式命中,会自动去重。远程路径(如s3://、http://)通过io_config参数配置访问凭据。
行与列的映射规则
每个样本由具有相同前缀的连续 TAR 成员构成,成员后缀成为列名。例如归档中依次包含000001.jpg、000001.json、000001.txt,则 Daft 生成一行,包含jpg、json、txt三列。此外 Daft 恒为每行附加两列元数据:
__key__:样本前缀(如000001);__url__:样本所在分片的路径。
该映射逻辑在源码_iter_archive_samples()中实现:解析成员名prefix.suffix(支持嵌套目录,nested/sample.en.txt会得到前缀nested/sample、字段en.txt,由测试 test_read_webdataset_preserves_multipart_suffixes 验证),遇到新前缀时产出上一个样本;同一前缀内若出现重复后缀(如两个a.txt)会直接报错,避免静默覆盖数据。
惰性媒体读取:只取元数据不读载荷
按后缀类型,成员被分为“立即解码”与“惰性引用”两类(对应源码中的_is_eagerly_decoded与_member_value):
- 立即解码:JSON(
json、jsn)、文本(txt、text、transcript)与整数类别/索引(cls、cls2、id、index、inx)。JSON 解析为 Python 对象、文本解码为 UTF-8 字符串、整数类转成int64; - 惰性引用:图片(
jpg、png、webp等 16 种)、音频(mp3、wav、flac等 11 种)、视频(mp4、mkv、mov等 10 种)及其它二进制成员,统一包装为带position(成员数据在 TAR 中的字节偏移)与size的daft.File引用,分别标注MediaType.image()、MediaType.audio()、MediaType.video()或unknown。
因此只选取元数据列不会触发媒体载荷的读取,这对超大规模数据集的快速浏览至关重要:
metadata = df.select("__key__", "json", "txt") images = df.select("__key__", "jpg")对images收集后,每行是一个ImageFile,仅在真正需要时通过.open()读取字节;远程分片(如 HTTP)则按Range请求精确拉取对应字节区间,测试 test_read_webdataset_http_range_references 验证了该行为。惰性引用机制同样支持列投影下推:源码WebDatasetSource.get_tasks()会根据Pushdowns.columns裁剪投影列,仅读取被选中的后缀成员。
从 Hugging Face 读取 WebDataset
Daft 的daft.read_huggingface()支持format="webdataset"参数,底层即转发到read_webdataset(f"hf://datasets/{repo}/**/*.tar"):
df = daft.read_huggingface( "laion/conceptual-captions-12m-webdataset", format="webdataset", )也可以直接传入显式的 Hugging Face 路径:
df = daft.read_webdataset( "hf://datasets/laion/conceptual-captions-12m-webdataset/data/*.tar" )format参数缺省时默认按"parquet"处理;传入其它值会抛出Unsupported Hugging Face dataset format错误。此外,read_webdataset的路径参数同样支持hf://协议,因此通过 URL 风格路径读取data/子目录下的所有分片即可绕过数据集根目录的自动递归。
格式要求与 Schema 推断策略
Daft 当前对 WebDataset 归档有以下硬性要求(源码与测试双重印证):
- 仅支持未压缩的
.tar:惰性成员读取依赖字节偏移(member.offset_data),压缩格式无法定位。源码_COMPRESSED_TAR_SUFFIXES明确列出.tar.gz、.tgz、.tar.bz2、.tar.xz等压缩后缀并直接拒绝,测试 test_read_webdataset_rejects_compressed_tar 验证报错Compressed WebDataset shards are not supported; - 不支持稀疏(sparse)TAR 成员:稀疏文件的逻辑内容不存放在连续字节区间,无法用偏移量引用,读取时报错并说明原因;
- 跨分片后缀集合与 JSON 结构必须一致:Schema 推断仅使用第一个分片的前 5 个样本(源码常量
_SAMPLES_FOR_SCHEMA_INFERENCE = 5)。若后续样本出现推断时未见过的后缀、或 JSON 字段结构与推断类型不符(如期望{caption, score}却收到{different}),Daft 会抛出ValueError而非静默丢弃数据。
该“宁可报错、绝不丢数据”的策略由三个测试直接验证:test_read_webdataset_rejects_field_missing_from_inferred_schema、test_read_webdataset_rejects_incompatible_json_schema 与 test_read_webdataset_rejects_duplicate_fields。
底层实现与执行模型
read_webdataset的完整调用链如下:
read_webdataset(path, io_config, batch_size)校验batch_size > 0(默认 1000,每次产出 RecordBatch 的最大样本数),未传io_config时回落到get_context().daft_planning_config.default_io_config;WebDatasetSource.__init__解析路径、对第一个分片执行_infer_schema推断出__key__、__url__及全部后缀列的类型;get_tasks()依据列投影与limit下推构造每个分片一个WebDatasetSourceTask(若设置了 limit,batch_size 取min(batch_size, limit));WebDatasetSourceTask.read()逐样本累积成 RecordBatch,通过asyncio.to_thread将 TAR 解析放到线程中执行,避免阻塞事件循环,再以异步迭代器产出。
从源码结构可以推断,WebDatasetSource继承自DataSource、任务继承DataSourceTask,与 Daft 其它连接器(Parquet、JSON 等)共用同一套分布式读取与下推框架,因此 WebDataset 读取天然支持列投影、limit 裁剪、本地/远程存储混用以及分布式执行。若样本缺失推断 Schema 中的某个字段,该字段取None(测试 test_read_webdataset_missing_inferred_field_is_null 验证),保证整表结构稳定。
常见错误速查
| 触发条件 | 报错信息 | 应对方式 |
|---|---|---|
传入.tar.gz/.tgz等压缩分片 | Compressed WebDataset shards are not supported | 先解压为纯.tar再读取 |
| 归档含稀疏成员 | Sparse TAR member ... is not supported | 重写归档,避免稀疏文件 |
| 后期样本出现新后缀 | was not present during schema inference | 统一所有分片的后缀集合,或将推断前置到包含全后缀的分片 |
| JSON 结构与推断不一致 | does not match its inferred schema | 统一各样本 JSON 的字段与类型 |
| 同一样本重复后缀 | Duplicate WebDataset field ... | 检查打包脚本,确保同前缀成员后缀唯一 |
batch_size <= 0 | batch_size must be greater than zero | 传正整数(默认 1000) |
未匹配到任何.tar | No uncompressed WebDataset TAR shards found | 检查路径、通配符与目录结构 |
小结
Daft 把 WebDataset 的 TAR 分片抽象为“一行一样本、后缀即列”的多模态表格,并以惰性File引用避免读取未选中的媒体字节,配合列投影下推、batch_size控制与分布式任务模型,可以高效完成海量图文/音视频数据集的元数据扫描、样本筛选与按需解码。使用前请务必遵守“未压缩 TAR、非稀疏成员、跨分片后缀与 JSON 结构一致”三项约束,即可将 Daft 作为 WebDataset 流水线的统一分析入口。
【免费下载链接】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),仅供参考