数据截至 (上游 commit a649de79c14a)
DataTrove — 架构与原理
30 秒导读: DataTrove 是 HuggingFace 的大规模文本数据处理库——FineWeb(15T token 的预训练语料)就是用它生产的。它把「从 Common Crawl 的 WARC 文件到训练用 token 二进制文件」这条长链路,拆成一串可自由拼装的
PipelineStep(读、抽取、过滤、去重、token 化、写出),每一步都是一个进出生成器的 Python 函数;同一条 pipeline 不改一行代码,就能在本机多进程或几千个 Slurm 任务上跑,且失败后重跑只补没完成的 shard。
1. 这是什么(零基础也能懂)
一句话定义
DataTrove 是一个预训练语料生产线框架:你给它一堆原始网页文件,它按你拼好的「流水线」逐文档清洗、去重、打分,最后输出干净的文本或 token 文件。
它要解决谁的什么问题
设想你要复现 FineWeb:手上有 Common Crawl 某次抓取的几十万个 .warc.gz 文件(里面是原始 HTML),想要一份能直接训练 LLM 的英文语料。你要做的一串事:
- 从 WARC 里抠出 HTML,再从 HTML 里抽出正文;
- 扔掉黄赌毒域名、非英语页、模板样板文本、复读机式重复内容;
- 在几十亿文档里找出近似重复(同一篇文章被反复转载)并只留一份;
- 最后 token 化、打乱文档顺序,写成训练框架能直接 mmap 读的二进制文件。
每件事单看都不难,难在规模:几十亿文档、TB 级数据,单机跑不动,必须切成几千份并行,还要能断点续跑、能统计「每一步扔掉了多少」。DataTrove 就是把这套「切分 + 拼装 + 续跑 + 统计」的基础设施一次性做好的那一层。
它能做什么
| 能力 | 具体支持 |
|---|---|
| 读取 | WARC(Common Crawl 原始格式)、JSONL、Parquet、CSV、Arrow IPC、HF Hub 数据集(src/datatrove/pipeline/readers/) |
| 正文抽取 | Trafilatura、readability 等 HTML → 正文抽取器(src/datatrove/pipeline/extractors/trafilatura.py) |
| 质量过滤 | Gopher / C4 / FineWeb 论文的启发式规则、URL 黑名单、fastText 语言识别与分类器、正则、采样(src/datatrove/pipeline/filters/) |
| 去重 | MinHash 近似去重(四阶段)、句子级精确去重(C4 风格)、整文精确去重、Bloom filter、ExactSubstr 子串去重(src/datatrove/pipeline/dedup/) |
| token 化 | HF tokenizers 批量编码成 .ds 二进制 + 文档边界索引 + loss mask;文档级/块级/窗口级混洗(src/datatrove/pipeline/tokens/) |
| 执行环境 | 本机多进程(LocalPipelineExecutor)、Slurm 集群(SlurmPipelineExecutor)、Ray、HF Jobs(src/datatrove/executor/) |
| 统计 | 每个 block 自动记录处理/丢弃数、文档长度分布、耗时,按任务落盘再合并(src/datatrove/utils/stats.py) |
用起来什么样
最小形态就是一个 Python 列表加一行 .run()(摘自 README.md:117-131,有精简):
from datatrove.executor.local import LocalPipelineExecutor
from datatrove.pipeline.readers import CSVReader
from datatrove.pipeline.filters import SamplerFilter
from datatrove.pipeline.writers import JsonlWriter
LocalPipelineExecutor(
pipeline=[
CSVReader(data_folder="/my/input/path"), # 读
SamplerFilter(rate=0.5), # 随机留一半
JsonlWriter(output_folder="/my/output/path"), # 写
],
tasks=10, # 切成 10 个 shard
workers=5, # 同时跑 5 个进程
).run()
把它放大到生产规模,就是 FineWeb 的真实配置(examples/fineweb.py:33-73):tasks=8000 个 Slurm 任务,每个任务读自己那份 WARC shard,依次过 URL 过滤 → Trafilatura 抽正文 → 语言识别 → Gopher 重复/质量 → C4 → FineWeb 质量过滤,每步被扔掉的文档由 exclusion_writer 另行存档(方便事后审计「这步扔了什么」),存活的写进 S3。
一句话直觉
把整条产线想成一根 Unix 管道:reader | filter1 | filter2 | writer。 每个环节都是 yield Document 的惰性生成器,数据流过即处理、不落地;真正特殊的是——这根管道会被竖切成 N 份(每份一个 rank,吃不同的文件 shard),并发跑在集群上,且每份跑完就敲一个「完成」章,崩了重跑只补没盖章的份。
2. 顶层全景(它大概怎么转)
2.1 顶层结构图
┌─────────────────────────────────────────┐
用户代码 │ Executor (Local / Slurm / Ray / Jobs) │
pipeline = [...] │ · 把 pipeline 按 rank 复制 N 份 │
│ │ · completions/ 空文件 = 断点续跑 │
▼ │ · stats/ 每任务统计 → 合并 │
└───────────────┬─────────────────────────┘
│ 每个 rank 内部:
▼
Reader ──► Filter ──► ... ──► Writer ──► Tokenizer
(取 rank (留/弃 + (模板文件 名 (.ds 二进制 +
的文件 exclusion_writer + 文件轮转 .index 边界 +
shard) 存被扔的) + HF 重试) 文档混洗)
│ │ ▲
└────────────┴── Document 生成器流 ┘
(text / id / metadata)
怎么读这张图: 上下两层。上层是 executor,唯一职责是「把同一条 pipeline 按 rank 跑 N 份并记账」;下层是单个 rank 内部的生成器链,Document 从左往右流,每个环节可以选择放行、修改、丢弃或旁路写盘。去重类步骤(MinHash 等)不在这条链上直接干活,而是拆成多个独立 job,靠落盘的签名文件串联——见第 4 章。
2.2 部件职责
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
Document | 流过管线的最小单元:text / id / metadata(+预留的 media) | src/datatrove/data.py:37-56 |
PipelineStep | 所有处理块的基类:依赖检查、stats、run(data, rank, world_size) | src/datatrove/pipeline/base.py:9 |
PipelineExecutor | executor 基类:逐 rank 调用 _run_for_rank、completions 记账、stats 落盘 | src/datatrove/executor/base.py:37 |
LocalPipelineExecutor | 本机多进程池跑任务 | src/datatrove/executor/local.py:15 |
SlurmPipelineExecutor | 把整个 executor 用 dill pickle 成 executor.pik,生成 sbatch 脚本提交 job array | src/datatrove/executor/slurm.py:48 |
DataFolder | fsspec 的目录封装:统一本地 / S3 / HF Hub 路径,负责列文件与切 shard | src/datatrove/io.py:89 |
BaseReader / BaseDiskReader | 读文件 → adapter → Document;按 rank 取文件 shard | src/datatrove/pipeline/readers/base.py:14、:113 |
DiskWriter | 写盘基类:文件名模板(${rank}/metadata 占位)、按大小轮转、HF 重试 | src/datatrove/pipeline/writers/disk_base.py:31 |
BaseFilter | 过滤器协议:`filter(doc) -> bool | (False, reason),可接 exclusion_writer` |
| MinHash 四件套 | MinhashDedupSignature / ...Buckets / ...Cluster / ...Filter 四阶段近似去重 | src/datatrove/pipeline/dedup/minhash.py:124、:324、:500、:599 |
DocumentTokenizer | 批量 token 化并写 .ds 二进制 + 索引,可选文档内混洗 | src/datatrove/pipeline/tokens/tokenizer.py:281 |
Stats / PipelineStats | 每个 block 的计数/均值/方差(Welford 在线算法 ),按任务落盘再并行合并 | src/datatrove/utils/stats.py:55、:132 |
2.3 主线走一遍(一个 rank 内部)
下面这条线对应 PipelineExecutor._run_for_rank(src/datatrove/executor/base.py:102)里的那段循环:
① 断点检查 is_rank_completed(rank) → completions/00042 存在就直接返回
(executor/base.py:116-118)
│
▼
② 串起生成器链 pipelined_data = None
for step in pipeline: pipelined_data = step(pipelined_data, rank, world_size)
(executor/base.py:129-136)
—— 注意:此刻只是「组装」,一 个文档都还没读
│
▼
③ 拉动数据 deque(pipelined_data, maxlen=0) (executor/base.py:137-138)
—— 惰性链条被拉到底:reader 开始读自己 shard 的文件
(DataFolder.get_shard: all_files[rank::world_size], io.py:164-180)
→ 每个 Document 依次流过 filter(可能被丢/被改) → writer(落盘)
│
▼
④ 记账 stats/{rank:05d}.json 落盘 + completions/{rank:05d} 空文件
(executor/base.py:143-148)
整条链是纯惰性的:② 只组装不执行,③ 的 deque(..., maxlen=0) 才是点火开关——它把生成器消费到底,数据才真正流动。这个设计让「N 个步骤的组合」和「单个步骤」有完全一样的接口,executor 不需要知道链上有什么。
跨 rank 的协作只发生在文件粒度:reader 用 all_files[rank::world_size] 这个隔行取方式切分输入文件,writer 用文件名里的 ${rank} 避免互相覆盖。而需要全局视野的步骤( 去重要看全数据集)走另一条路:拆成多个 job,中间结果落盘(签名文件 → 重复对 → 删除清单),用 executor 的 depends= 串成 DAG——FineWeb 的 MinHash 就是这样四个 Slurm job 串起来的(examples/fineweb.py:103-180)。
3. 阅读地图(建议顺序)
五章由浅入深。时间有限就读 01 → 04:01 是整套系统的骨架,04 是这个库最有工程含量的部分。
| 顺序 | 章节 | 讲什么 | 适合谁 |
|---|---|---|---|
| 1 | 01-executors-pipeline.md | PipelineStep 生成器链、分片、completions 续跑、三种 executor | 所有人必读,骨架在这 |
| 2 | 02-io-readers-writers.md | DataFolder 统一本地/远端路径、adapter、文件名模板、文件轮转 | 要接入自己的数据格式的人 |
| 3 | 03-filters.md | Gopher/C4/FineWeb 的规则原文与实现、URL/语言/分类器过滤 | 关心「好语料长什么样」的人 |
| 4 | 04-dedup.md | MinHash 四阶段、句子级精确去重、Bloom filter;分布式 k 路归并与并查集 | 想做亿级去重的人 |
| 5 | 05-tokens.md | .ds 二进制格式、loss mask、三级混洗、跨文件 merger | 要产出训练用 token 文件的人 |
4. 巧妙之处(可借鉴的技术)
每一条的细节和源码引用都在对应章节,这里先给「精华预览」:
- 生成器链 +
deque(..., maxlen=0)点火(src/datatrove/executor/base.py:137-138):整条 pipeline 零缓存、逐文档流动,单进程内存占用与数据集大小无关——这是它能用mem_per_cpu_gb=2的小任务跑完整个 Common Crawl 的根本原因(examples/fineweb.py:70)。 - completions 空文件 = 最简断点续跑(
src/datatrove/executor/base.py:156-177):做完一个 rank 就写一个空文件,重跑时get_incomplete_ranks列目录差集即可。Slurm 上重投同一个 job 就自动只跑没完成的任务。 - 被丢弃的数据也是一等公民:每个过滤器都能挂
exclusion_writer把扔掉的文档连原因一起存盘(src/datatrove/pipeline/filters/base_filter.py:61-82),「哪步扔了 30%」随时可以审计。 - 去重的多阶段落盘设计:MinHash 不在内存里做全局比对,而是「签名排序落盘 → 按桶分片做 k 路堆归并 → 并查集聚类 → 按
.remove清单过滤」四个 job(src/datatrove/pipeline/dedup/minhash.py)。每阶段都是纯顺序 I/O,能扛任意大数据量;签名文件排序后还会用 blake2b 校验写盘正确性(minhash.py:290-321)——注释里明说「未排序的签名是 MinHash 的头号事故源」。 - 文件名模板
${tag}直接吃 metadata(src/datatrove/pipeline/writers/disk_base.py:166-185):output_filename="${language}/dump/${rank}.jsonl.gz"就实现了按语言分目录归档,零额外代码。 - 在线统计可并行合并:
MetricStats用 Welford 在线算法算均值/方差,__add__用并行方差公式把几千个任务的 stats 合成全局报告(src/datatrove/utils/stats.py:217-274)。 - 依赖惰性检查:
PipelineStep.__new__在实例化时按_requires_dependencies检查并提示装什么(src/datatrove/pipeline/base.py:22-32),用户不必为用不到的 block 装重依赖(fasttext、spacy 等)。
5. 边界与局限
诚实地说,这个库不做什么、在哪会疼:
- 文件是切分的最小单位。 一个文件只会被一个 task 完整处理,不会自动拆分(
README.md:85-86的 TIP 明说)。输入是少量超大文件时,并行度被文件数卡死。 - 阈值是经验值,不是原理。 Gopher/C4 的阈值照抄论文消融实验(
gopher_repetition_filter.py:11-28直接把论文 Table A1 抄成注释);换个语言/领域要自己重调。词/句切分按语言选 tokenizer,没覆盖的语言要自己加(src/datatrove/utils/word_tokenizers.py:59-73)。 - 部分步骤是天然单点的。 MinHash stage 3(并查集聚类)必须
world_size == 1(minhash.py:533)且内存要装得下全 部重复对——FineWeb 给它mem_per_cpu_gb=25, cpus_per_task=8(examples/fineweb.py:154-156)。DocumentTokenizerMerger同理只能单任务跑(src/datatrove/pipeline/tokens/merger.py:100)。 - 去重各阶段靠文件约定串联。 stage 2 要求
world_size % num_buckets == 0(minhash.py:393);sentence dedup 的 finder 任务数必须等于 stage 1 的finder_workers,否则直接报错(sentence_dedup.py:222-228)。配错不会静默错,但要读懂报错得理解这套分桶约定。 - 为集群设计,本地运行只是调试档。 它的心智模型是「几千个 task 分片跑」,本地 executor 适合小样例;
shuffle_files这类选项在 dedup 块前使用会破坏 doc 顺序假设(readers/base.py:132的 docstring 警告)。 - 远端混洗/合并很慢。 tokenizer 的混洗靠随机 seek 读(
tokenizer.py:15-17特意关了缓存、只读 50KB 块),merger 的 docstring 直接警告 S3 上会慢(merger.py:19-21)——中间结果尽量落本地盘。
6. 横向对比
同书架上处理同类问题的项目,取舍各不相同:
- Dolma(AI2 的语料工具链):两者都是「过滤 + 去重」产线,目标几乎一致。Dolma 把规则组织成 tagger 生态并以 Rust 重活加速;DataTrove 保持纯 Python + 生成器,用「多阶段落盘 + executor DAG」扛规模,且自带 Slurm/HF 生态集成。FineWeb 论文(Penedo et al., 2024)里 DataTrove 侧的主张是「轻依赖、 少抽象」。
- hf-datasets:HF Datasets 解决的是「一份整理好的数据集怎么高效存储/流式读取/版本化」,面向使用者;DataTrove 解决的是「怎么从原始网页把数据集造出来」,面向生产者。两者通过
hf://路径和 HF writer 衔接(DataTrove 产物常发布为 datasets)。 - OLMo:OLMo 是完整训练栈,语料处理只是其中一环(且主要用 Dolma);DataTrove 专注且只做数据产线,不含任何训练代码。
7. 代码地图(导航索引)
从零开始读源码的推荐入口,按主题跳:
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 数据单元 | src/datatrove/data.py | Document、DocumentsPipeline |
| 处理块基类 | src/datatrove/pipeline/base.py | PipelineStep、PipelineStep.stat_update |
| executor 骨架(串链 + 续跑) | src/datatrove/executor/base.py | PipelineExecutor._run_for_rank、is_rank_completed |
| 本地执行 | src/datatrove/executor/local.py | LocalPipelineExecutor.run |
| Slurm 执行(pickle + sbatch) | src/datatrove/executor/slurm.py | SlurmPipelineExecutor.launch_job、launch_slurm_job |
| 路径/文件系统抽象 | src/datatrove/io.py | DataFolder.get_shard、get_datafolder、safely_create_file |
| 读取协议 | src/datatrove/pipeline/readers/base.py | BaseReader._default_adapter、BaseDiskReader.run |
| 写出协议 | src/datatrove/pipeline/writers/disk_base.py | DiskWriter.write、_get_output_filename |
| 过滤器协议 | src/datatrove/pipeline/filters/base_filter.py | BaseFilter.filter、BaseFilter.run |
| 启发式质量过滤 | src/datatrove/pipeline/filters/gopher_quality_filter.py | GopherQualityFilter.filter |
| 重复度过滤 | src/datatrove/pipeline/filters/gopher_repetition_filter.py | GopherRepetitionFilter、find_all_duplicate |
| MinHash 去重(四阶段) | src/datatrove/pipeline/dedup/minhash.py | MinhashDedupSignature、MinhashDedupBuckets、MinhashDedupCluster、MinhashDedupFilter |
| 句子级去重 | src/datatrove/pipeline/dedup/sentence_dedup.py | SentenceDedupSignature、SentenceDedupFilter.remove_dup_sentences |
| 文本归一化 | src/datatrove/utils/text.py | simplify_text、ngrams |
| 哈希配置 | src/datatrove/utils/hashing.py | HashConfig、create_hash_func |
| token 化写出 | src/datatrove/pipeline/tokens/tokenizer.py | DocumentTokenizer、TokenizedFile |
| 跨文件合并混洗 | src/datatrove/pipeline/tokens/merger.py | DocumentTokenizerMerger、load_doc_ends |
| 在线统计 | src/datatrove/utils/stats.py | MetricStats、PipelineStats |
| FineWeb 生产配置 | examples/fineweb.py | main_processing_executor、stage1~stage4 |