跳到主要内容

数据截至 (上游 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(+预留的 mediasrc/datatrove/data.py:37-56
PipelineStep所有处理块的基类:依赖检查、stats、run(data, rank, world_size)src/datatrove/pipeline/base.py:9
PipelineExecutorexecutor 基类:逐 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 arraysrc/datatrove/executor/slurm.py:48
DataFolderfsspec 的目录封装:统一本地 / S3 / HF Hub 路径,负责列文件与切 shardsrc/datatrove/io.py:89
BaseReader / BaseDiskReader读文件 → adapter → Document;按 rank 取文件 shardsrc/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_ranksrc/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 是这个库最有工程含量的部分。

顺序章节讲什么适合谁
101-executors-pipeline.mdPipelineStep 生成器链、分片、completions 续跑、三种 executor所有人必读,骨架在这
202-io-readers-writers.mdDataFolder 统一本地/远端路径、adapter、文件名模板、文件轮转要接入自己的数据格式的人
303-filters.mdGopher/C4/FineWeb 的规则原文与实现、URL/语言/分类器过滤关心「好语料长什么样」的人
404-dedup.mdMinHash 四阶段、句子级精确去重、Bloom filter;分布式 k 路归并与并查集想做亿级去重的人
505-tokens.md.ds 二进制格式、loss mask、三级混洗、跨文件 merger要产出训练用 token 文件的人

4. 巧妙之处(可借鉴的技术)

每一条的细节和源码引用都在对应章节,这里先给「精华预览」:

  1. 生成器链 + deque(..., maxlen=0) 点火src/datatrove/executor/base.py:137-138):整条 pipeline 零缓存、逐文档流动,单进程内存占用与数据集大小无关——这是它能用 mem_per_cpu_gb=2 的小任务跑完整个 Common Crawl 的根本原因(examples/fineweb.py:70)。
  2. completions 空文件 = 最简断点续跑src/datatrove/executor/base.py:156-177):做完一个 rank 就写一个空文件,重跑时 get_incomplete_ranks 列目录差集即可。Slurm 上重投同一个 job 就自动只跑没完成的任务。
  3. 被丢弃的数据也是一等公民:每个过滤器都能挂 exclusion_writer 把扔掉的文档连原因一起存盘(src/datatrove/pipeline/filters/base_filter.py:61-82),「哪步扔了 30%」随时可以审计。
  4. 去重的多阶段落盘设计:MinHash 不在内存里做全局比对,而是「签名排序落盘 → 按桶分片做 k 路堆归并 → 并查集聚类 → 按 .remove 清单过滤」四个 job(src/datatrove/pipeline/dedup/minhash.py)。每阶段都是纯顺序 I/O,能扛任意大数据量;签名文件排序后还会用 blake2b 校验写盘正确性(minhash.py:290-321)——注释里明说「未排序的签名是 MinHash 的头号事故源」。
  5. 文件名模板 ${tag} 直接吃 metadatasrc/datatrove/pipeline/writers/disk_base.py:166-185):output_filename="${language}/dump/${rank}.jsonl.gz" 就实现了按语言分目录归档,零额外代码。
  6. 在线统计可并行合并MetricStats 用 Welford 在线算法算均值/方差,__add__ 用并行方差公式把几千个任务的 stats 合成全局报告(src/datatrove/utils/stats.py:217-274)。
  7. 依赖惰性检查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 == 1minhash.py:533)且内存要装得下全部重复对——FineWeb 给它 mem_per_cpu_gb=25, cpus_per_task=8examples/fineweb.py:154-156)。DocumentTokenizerMerger 同理只能单任务跑(src/datatrove/pipeline/tokens/merger.py:100)。
  • 去重各阶段靠文件约定串联。 stage 2 要求 world_size % num_buckets == 0minhash.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.pyDocumentDocumentsPipeline
处理块基类src/datatrove/pipeline/base.pyPipelineStepPipelineStep.stat_update
executor 骨架(串链 + 续跑)src/datatrove/executor/base.pyPipelineExecutor._run_for_rankis_rank_completed
本地执行src/datatrove/executor/local.pyLocalPipelineExecutor.run
Slurm 执行(pickle + sbatch)src/datatrove/executor/slurm.pySlurmPipelineExecutor.launch_joblaunch_slurm_job
路径/文件系统抽象src/datatrove/io.pyDataFolder.get_shardget_datafoldersafely_create_file
读取协议src/datatrove/pipeline/readers/base.pyBaseReader._default_adapterBaseDiskReader.run
写出协议src/datatrove/pipeline/writers/disk_base.pyDiskWriter.write_get_output_filename
过滤器协议src/datatrove/pipeline/filters/base_filter.pyBaseFilter.filterBaseFilter.run
启发式质量过滤src/datatrove/pipeline/filters/gopher_quality_filter.pyGopherQualityFilter.filter
重复度过滤src/datatrove/pipeline/filters/gopher_repetition_filter.pyGopherRepetitionFilterfind_all_duplicate
MinHash 去重(四阶段)src/datatrove/pipeline/dedup/minhash.pyMinhashDedupSignatureMinhashDedupBucketsMinhashDedupClusterMinhashDedupFilter
句子级去重src/datatrove/pipeline/dedup/sentence_dedup.pySentenceDedupSignatureSentenceDedupFilter.remove_dup_sentences
文本归一化src/datatrove/utils/text.pysimplify_textngrams
哈希配置src/datatrove/utils/hashing.pyHashConfigcreate_hash_func
token 化写出src/datatrove/pipeline/tokens/tokenizer.pyDocumentTokenizerTokenizedFile
跨文件合并混洗src/datatrove/pipeline/tokens/merger.pyDocumentTokenizerMergerload_doc_ends
在线统计src/datatrove/utils/stats.pyMetricStatsPipelineStats
FineWeb 生产配置examples/fineweb.pymain_processing_executorstage1~stage4