数据截至 (上游 commit a649de79c14a)
02 · IO 层、读取器与写出器
这一章讲什么: 数据进出的两端。读侧:
DataFolder如何用 fsspec 统一本地/S3/HF 路径并按 rank 切文件 shard,reader 如何用 adapter 把一行 JSON/一条 WARC 记录变成Document。写侧:DiskWriter的文件名模板、按大小轮转,以及为 HF Hub 专门写的重试逻辑。
1. 它要解决的小 问题
产线的两端各有一坨琐碎但致命的工程问题:
- 读侧:数据可能在本地盘、S3 bucket 或 HF Hub 上;一个 job 的几千个 rank 要确定性地瓜分文件列表,不能重叠也不能漏;输入格式五花八门(WARC、JSONL、Parquet……),但下游只认
Document。 - 写侧:几千个 rank 同时写,文件名不能撞;单个文件不能无限大;写到 HF Hub 会被限流(429)或撞上 commit 竞态(412),不能让整个任务因此挂掉。
这一层的答案是两个基类加一个路径封装:BaseDiskReader 管「列文件 → 切 shard → adapter → Document」,DiskWriter 管「文件名模板 → 轮转 → 容错写」,DataFolder 管「这句话在哪个文件系统上执行」。
2. DataFolder:会切 shard 的目录
DataFolder(src/datatrove/io.py:89)继承 fsspec 的 DirFileSystem,构造时用 url_to_fs 从路径字符串推断文件系统(:101-116)——s3://bucket/path 得到 S3,hf://datasets/... 得到 HF Hub,裸路径得到本地盘。用户也可以直接传 (str, fs对象) 或 (str, dict),由 get_datafolder 工厂统一消化(:251-283)。
它对产线最重要的方法是 get_shard(:164-180):
# src/datatrove/io.py:177-180
all_files = self.list_files(**kwargs)
if len(all_files) == 0:
return None
return all_files[rank::world_size]
隔 world_size 取一份——列表是排序过的(list_files 返回 sorted(...),:147-162),所以这个划分是确定性的:只要文件集合和 world_size 不变,每个 rank 永远拿到同样的文件,重跑不会产生重叠。这就是第 1 章「tasks 数不能改」的原因。
两个配套小件:
OutputFileManager(:17-86):writer 的文件句柄池,按文件名缓存打开的文件、统一关闭;safely_create_file(:353-379):用fasteners.InterProcessLock+.completed标记文件实现「多进程安全地下载/生成某个资源一次」。需要下载模型/黑名单的 block(语言识别、URL 过滤)都靠它,封装成cached_asset_path_or_download(:382-405)。
3. 读侧:adapter 协议
3.1 直觉
所有输入格式的差异被 压缩成一个函数:给我一个 dict,还我 Document 的字段。
# 示意,非源码
def my_adapter(self, data: dict, path: str, id_in_file: int):
return {
"text": data.pop("content"), # 必须有 text
"id": data.pop("url"), # 必须有 id
"metadata": {"source": path} | data, # 剩下的塞进 metadata
}
默认 adapter(BaseReader._default_adapter,src/datatrove/pipeline/readers/base.py:49-76)做的事就是:按 text_key/id_key 取字段(取不到 id 就用 f"{path}/{id_in_file}" 兜底),metadata 字段若是字符串尝试 JSON 解析,其余字段统统并入 metadata。所以一个上游有 {"content": ..., "url": ..., "date": ...} 的数据集,只需 text_key="content", id_key="url" 两个参数就接上了。
3.2 真实实现
类层级只有两层:BaseReader(readers/base.py:14)管 adapter 与 Document 构建;BaseDiskReader(:113)管「文件 → 行」的骨架。子类只需实现一个方法:
read_file(filepath):打开文件,逐条记录 yield(抽象在:171-182)。
骨架里的关键流转在 run(:224-252)和 read_files_shard(:184-222):前者先透传上游数据(if data: yield from data,:235-236——所以 reader 也能接在别的块后面),再用 get_shard 取本 rank 的文件;后者逐文件逐文档 yield,顺手做 skip/limit 截断和进度条。
get_document_from_dict(:78-103)里有个防御细节:adapter 产出的 text 为空时直接丢弃该文档,但警告只打一次(_empty_warning),并提示「是不是 text_key 配错了,现有键是这些」——大规模跑时这是最高频的配置错误。
3.3 两个典型 reader
JsonlReader(src/datatrove/pipeline/readers/jsonl.py:9):逐行 orjson 解析;单行坏了只警告不炸(JSONDecodeError时continue,:81-91),整文件编码坏了降级为警告(:93-94)——几十亿行里必然有坏行,不能为它停任务。WarcReader(src/datatrove/pipeline/readers/warc.py:11):Common Crawl 的入口。process_record(:87-140)只保留response/conversion记录、只收text/html等少数 MIME(老的 dump 没有 MIME 头就用python-magic嗅探,:106-114),编码先按 UTF-8 解、失败用 cchardet 探测再试一次(:116-129),最后把WARC-Record-ID/WARC-Target-URI/WARC-Date塞进 dict——URL 就是这样进入 metadata、供 URL 过滤和统计分域使用的。
4. 写侧:文件名模板与容错
4.1 直觉
写出器要回答的问题不是「怎么写」,而是「写进哪个文件」。DataTrove 把答案做成一个 string.Template:占位符当场从文档身上取。
${rank}—— 当前任务的 rank(补零到 5 位);${id}—— 文档 id;${任意词}—— 文档metadata里的同名键。
于是 output_filename="${language}/" + DUMP + "/${rank}.jsonl.gz" 一行就实现了「按语言分目录、按 rank 分文件」(FineWeb 里就是这么归档非英语文档的,examples/fineweb.py:46)。
4.2 真实实现
模板替换在 DiskWriter._get_output_filename(src/datatrove/pipeline/writers/disk_base.py:166-185):
# src/datatrove/pipeline/writers/disk_base.py:183-185
return self.output_filename.substitute(
{"rank": str(rank).zfill(5), "id": document.id, **document.metadata, **kwargs}
)
构造期还有个贴心警告:如果模板里没有 ${rank},会警告「并发 worker 可能互相覆盖输出」(:124-128)。
write(:281-307)里的第二件事是按大小轮转:max_file_size > 0 时文件名前加 000_/001_ 计数器,当前文件 .tell() 超阈值就 _on_file_switch 关旧开新(:294-302)。注意它要求二进制模式(文本模式的 tell 不可信),构造期直接 ValueError(:121-122)。
第三件事是 adapter 的镜像:_default_adapter(:134-155)把 Document 转回 dict,expand_metadata=True 时把 metadata 摊平成顶层列(写给 Parquet 这类列式格式用)。
4.3 HF Hub 重试
往 HF Hub 写是「先写本地临时文件、close 时上传」的模式,在几千并发 writer 下会稳定遇到 429/412/503。DiskWriter 对此有专门一层(disk_base.py:48-98):
HF_RETRYABLE_MESSAGES(:18-28)列了一串可 重试的错误子串(限流、commit 竞态、LFS 校验失败等);_retry_hf_hub_operation(:66-86)做指数退避,最多 12 次;- 最妙的是
close_file(:212-266):fsspec 在 close 上传失败后会把文件标记为已关闭,但临时文件还在盘上——于是它直接拿file_obj.fs._api手动重试upload_file,bucket 路径和 repo 路径走不同 API。注释里明说这个分支就是为了兜住这个状态不一致(:233-236)。
5. 关键细节 / 坑
- 文件名模板用
substitute而非safe_substitute:metadata 里缺键会直接KeyError(_get_output_filename用的是substitute,:183)。用 metadata 占位前确认上游一定写了这个键(比如LanguageFilter会写language,filters/language_filter.py:55)。 - reader 不是只能当链首:
BaseDiskReader.run开头透传上游数据(readers/base.py:235-236),reader 也可以出现在链中间做「边读边并流」;但要意识到此时两个来源的文档混在一起往下游走。 shuffle_files=True与 dedup 块不兼容(readers/base.py:132docstring):去重依赖「rank 内文档顺序稳定」这个隐性契约。- 压缩自动推断:writer 的
compression="infer"从文件名后缀猜;显式给gzip/zstd而文件名没带对应后 缀时会自动补后缀(disk_base.py:115-118)——输出文件名未必等于你写的模板。 - 代码里看不出的:各远端文件系统的实际吞吐/并发上限取决于 fsspec 实现和集群网络,库本身只提供
upload_block_size这类透传旋钮(如DocumentTokenizer的 S3 分块上传提示,tokens/tokenizer.py:325-326)。
6. 本章代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 目录与 shard | src/datatrove/io.py | DataFolder.get_shard、list_files、get_datafolder |
| 句柄池 | src/datatrove/io.py | OutputFileManager |
| 进程安全的资源下载 | src/datatrove/io.py | safely_create_file、cached_asset_path_or_download |
| reader 基类 | src/datatrove/pipeline/readers/base.py | BaseReader._default_adapter、BaseDiskReader.run、read_files_shard |
| JSONL 读 | src/datatrove/pipeline/readers/jsonl.py | JsonlReader.read_file |
| WARC 读 | src/datatrove/pipeline/readers/warc.py | WarcReader、process_record |
| writer 基类 | src/datatrove/pipeline/writers/disk_base.py | DiskWriter.write、_get_output_filename、close_file、_retry_hf_hub_operation |
| JSONL 写 | src/datatrove/pipeline/writers/jsonl.py | JsonlWriter._write |