跳到主要内容

数据截至 (上游 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 的目录

DataFoldersrc/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_adaptersrc/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 真实实现

类层级只有两层:BaseReaderreaders/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

  • JsonlReadersrc/datatrove/pipeline/readers/jsonl.py:9):逐行 orjson 解析;单行坏了只警告不炸JSONDecodeErrorcontinue:81-91),整文件编码坏了降级为警告(:93-94)——几十亿行里必然有坏行,不能为它停任务。
  • WarcReadersrc/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_filenamesrc/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 会写 languagefilters/language_filter.py:55)。
  • reader 不是只能当链首BaseDiskReader.run 开头透传上游数据(readers/base.py:235-236),reader 也可以出现在链中间做「边读边并流」;但要意识到此时两个来源的文档混在一起往下游走。
  • shuffle_files=True 与 dedup 块不兼容readers/base.py:132 docstring):去重依赖「rank 内文档顺序稳定」这个隐性契约。
  • 压缩自动推断:writer 的 compression="infer" 从文件名后缀猜;显式给 gzip/zstd 而文件名没带对应后缀时会自动补后缀(disk_base.py:115-118)——输出文件名未必等于你写的模板。
  • 代码里看不出的:各远端文件系统的实际吞吐/并发上限取决于 fsspec 实现和集群网络,库本身只提供 upload_block_size 这类透传旋钮(如 DocumentTokenizer 的 S3 分块上传提示,tokens/tokenizer.py:325-326)。

6. 本章代码地图

主题文件路径符号名
目录与 shardsrc/datatrove/io.pyDataFolder.get_shardlist_filesget_datafolder
句柄池src/datatrove/io.pyOutputFileManager
进程安全的资源下载src/datatrove/io.pysafely_create_filecached_asset_path_or_download
reader 基类src/datatrove/pipeline/readers/base.pyBaseReader._default_adapterBaseDiskReader.runread_files_shard
JSONL 读src/datatrove/pipeline/readers/jsonl.pyJsonlReader.read_file
WARC 读src/datatrove/pipeline/readers/warc.pyWarcReaderprocess_record
writer 基类src/datatrove/pipeline/writers/disk_base.pyDiskWriter.write_get_output_filenameclose_file_retry_hf_hub_operation
JSONL 写src/datatrove/pipeline/writers/jsonl.pyJsonlWriter._write