跳到主要内容

数据截至 (上游 commit 669f534823b0)

Dolma — 架构与原理

30 秒导读: Dolma 是 AI2 为训练 OLMo 而开源的两样东西:一个 3 万亿 token 的开放预训练语料,以及构建它的工具箱(本仓库)。工具箱的核心思想一句话:先给每篇文档打各种「属性标签」并存到旁边的 attributes 文件里,最后用一个 mixer 按配置决定留谁、删哪段、替换什么——过滤规则是数据,不是代码。


1. 这是什么(零基础也能懂)

一句话定义

Dolma Toolkit 是一个大规模语料清洗/构建工具箱:输入是几十亿篇 JSONL 文档(网页、论文、代码、书籍、Reddit、维基),输出是一个去重、去污、按比例混合好、可以直接 tokenize 进训练集的语料。

它要解决谁的什么问题

假设你要从零训一个开源大模型,第一步不是写模型代码,而是回答:训练数据从哪来、怎么洗?

  • Common Crawl 里有几十亿个网页,但大量是重复内容、机器生成垃圾、带个人隐私(邮箱/电话/IP)。
  • 评测集的题目可能泄漏在网页里(decontamination,去污),不洗掉分数就虚高。
  • 每种清洗规则(去重、语言识别、毒性过滤)都要扫一遍全量数据——扫一遍就是几 TB 的 IO,扫十遍谁都受不了。
  • 洗完之后还要按比例混:网页占多少、代码占多少、论文占多少,这个「配方」你会想反复调,不想每调一次就重洗一遍。

Dolma 就是 AI2 回答这些问题的工程成果,也是 OLMo 系列模型的数据侧。

它能做什么

能力命令 / 部件在哪
打标签(质量/语言/毒性/PII/长度等 30+ 种)dolma tag + tagger 注册表python/dolma/taggers/
去重(URL 级、段落级、ngram 重叠级)dolma dedupe + Rust Bloom filtersrc/deduper.rssrc/bloom_filter.rs
按配方混合 + 过滤 + span 替换dolma mixsrc/mixer.rssrc/shard.rs
tokenize 成训练用 memmapdolma tokenspython/dolma/tokenizer/
WARC(Common Crawl 原始格式)处理dolma warcpython/dolma/warc/
属性统计/可视化分析dolma analyzepython/dolma/core/analyzer.py

用起来什么样

一条典型的生产流水线是三条命令(摘自 docs/taggers.mddocs/deduplication.md 与真实配置 configs/dolma-v1_6/mixing/c4.yaml 的用法):

# ① 打标:对文档跑 gopher 规则,结果写到 .../attributes/gopher_v1/
dolma tag --documents 's3://bucket/cc/documents/*.json.gz' \
--taggers gopher_v1 --processes 32

# ② 去重:段落级,结果写到 .../attributes/dedupe_paragraphs/
dolma -c dedupe-paragraphs.json dedupe

# ③ 混合:按 YAML 配方过滤 + 替换 span,产出最终 documents/
dolma -c mixing/c4.yaml mix

注意 ①② 都不改动原文——它们只在平行的 attributes/ 目录里写标注文件;真正「动刀」的只有 ③。

一句话直觉

把 Dolma 想成医院的「检查报告 → 会诊」模式。 每个 tagger 是一项检查(验血、拍片),报告(attributes)一张张夹在病历里,谁都不许直接给病人动手术;最后 mixer 这个「主治医生」拿着全部报告,按一张治疗方案(混合配置)一次性决定:哪个病人出院、哪个要切掉哪一块。


2. 顶层全景(它大概怎么转)

2.1 目录约定:整个系统的地基

Dolma 的一切机制都建立在一条路径约定上(见 docs/data-format.md):

dataset/
├── documents/ ← 原始文档,一行一篇 {id, text, source, metadata}
│ └── 2019-09/0000.jsonl.gz
└── attributes/ ← 标注,与 documents 镜像同构、逐行对齐
├── gopher_v1/2019-09/0000.jsonl.gz ← {id, attributes: {…}}
├── dedupe_paragraphs/2019-09/0000.jsonl.gz
└── pii_detection/2019-09/0000.jsonl.gz

怎么读: attributes/<名字>/ 下的每个文件,与 documents/ 下同名文件行数相同、顺序相同、id 逐行对应。替换路径里的 /documents//attributes/<名>/ 就完成寻址——这个字符串替换硬编码在 Rust 侧 src/shard.rs:53(split_streams 内的 input.replace("/documents/", &attr_prefix))。

2.2 部件与数据流

┌──────────── Python 侧 ────────────┐
│ dolma tag (tagger 注册表 + │
│ BaseParallelProcessor 多进程) │
s3://…/documents ──► 每个 tagger 扫一遍 ──► s3://…/attributes/<exp>/
(jsonl.gz) │
└──────────────────────────────┼────┘
┌──────────── Rust 侧 ─────────┼────┐
│ dolma dedupe (Bloom filter) │ │
documents ────► 段落/URL 哈希查重 ───────────►├─► attributes/dedupe_*/
│ ▼ │
│ dolma mix (shard 装箱) │
documents + attributes/* ──► 按行合并 ──► 过滤/替换 ──► 最终 documents/
└───────────────────────────────────┘

┌──────────── Python 侧 ────────────┐
│ dolma tokens (ring 混洗 + memmap) │
最终 documents ──► tokenize ──► part-*.npy + .csv.gz │
└───────────────────────────────────┘

怎么读这张图: 自上而下是时间顺序。前两阶段只产「报告」,第三阶段「会诊动刀」,第四阶段把成品转成训练框架直接 mmap 读取的 token 流。Python 负责灵活(tagger 生态),Rust 负责重活(去重与混合的热路径),两者通过 pyo3 桥接(src/lib.rs:114-115#[pymodule] fn dolma)。

2.3 部件职责

部件干什么在哪个文件
BaseParallelProcessor「一个文件 = 一个任务」的多进程框架,带断点续跑(.done.txt)与进度条python/dolma/core/parallel.py:52
TaggerRegistry / BaseTaggertagger 注册表与基类;predict 返回 span 列表python/dolma/core/registry.py:9python/dolma/core/taggers.py:25
TaggerProcessor一遍扫描内跑多个 tagger,按实验名分文件写出python/dolma/core/runtime.py:225
BloomFilter(Rust)线程安全位数组 + k 个独立种子的 ahash;可落盘/读回src/bloom_filter.rs:15
deduper::run文档/段落/ngram 三级去重,写 span 属性src/deduper.rs:21
mixer::run + Shard按大小装箱、合并属性、过滤、span 替换、写最终文件src/mixer.rs:11src/shard.rs:17
DocFilter / SpanReplacerJSONPath 与 jq 双语法的「留/弃」判定与文本替换src/filters.rs:452src/shard.rs:599
MemMapParallelWriter多进程 tokenize + 环形缓冲混洗 + .npy 落盘python/dolma/tokenizer/executor.py:30

2.4 主线走一遍(以 C4 为例)

configs/dolma-v1_6/mixing/c4.yaml 是真实的生产配方,它假定下游已经发生过:

  1. 打标:gopher_v1 把 20 个质量信号(词数、词长中位数、重复 ngram 占比……)写成 span 属性。
  2. 去重:dedupe_docs(URL 级)与 dedupe_paragraphs(段落级)各自写一个属性文件。
  3. PII:pii_regex_with_counts_fast_v2 标出邮箱/电话/IP 的位置与数量。

然后 dolma mix 一遍完成:按行把 4 份 attributes 与原文合并 → 用 20 余条 JSONPath exclude 规则丢文档(如「词数 < 50」「PII 超过 5 处」)→ 对留下的文档做 span 替换(重复段落删掉、邮箱换成 |||EMAIL_ADDRESS|||)→ 按 4GB 装箱写出,并给每行打上 文件名:行号 的 provenance(src/shard.rs:455)。


3. 阅读地图

按由浅入深的顺序,建议这样读:

顺序章节你会学到
101-toolkit-architecture.md为什么 documents/attributes 要分开;并行框架怎么用「.done.txt」白嫖断点续跑;tagger 输出为什么长 exp__tagger__type 这样
202-dedup-bloom-filter.mdBloom filter 的尺寸/哈希数怎么算;为什么去重可以「天然并发」;怎么把 key 按哈希分区实现多机去重
303-mixer-span-surgery.mdmixer 一遍 IO 完成合并+过滤+替换的全过程;span 替换那段逐字符状态机;jq 与 JSONPath 两种语法的差异
404-taggers-catalog.md内置 tagger 的分类学:启发式规则(Gopher/C4)、正则(PII)、外部黑名单(URL/adblock)、模型(fastText/语言)
505-tokenize-memmap.md为什么「ring buffer + 局部 shuffle」就够混;memmap 预分配再裁尾的写法

只想知道「这个库值不值得抄」:读 index 的 §5(巧妙之处)与各章末尾的「坑」。


4. 核心机制速览(章节的种子)

  • 标注与决策分离:tagger 永远不删数据,只产 [start, end, score] 三元组;删不删是 mixer 配置里的一条表达式。这让「调配方」变成改 YAML 而不是重跑标注。详见 01、03 章。
  • Bloom filter 即全局去重状态:一个 Vec<AtomicU32> 位数组 + 若干独立种子的 ahash,contains 为假才 insert,所以多线程并发只会带来少量额外假阳性,不会漏放重复。详见 02 章。
  • 一切过滤都是表达式:mixer 的 include/exclude 是 JSONPath 或 jq 表达式,直接写在 YAML 里,对着合并后的文档逐行求值。详见 03 章。
  • 混洗是「够用就好」:tokenizer 不做全局 shuffle,而是 8 个文件轮转读、每 1 万条局部洗牌。代价近似零,效果接近全局混洗。详见 05 章。

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

  1. .done.txt 断点续跑:每处理完一个文件,在 metadata 目录写一个同名 .done.txt 时间戳文件;下次启动直接跳过已完成的(python/dolma/core/parallel.py:481existing_metadata_names 检查与 python/dolma/core/parallel.py:229 的写回)。成本几乎为零,却能让跑了几天的任务在机器挂了之后原地复活。
  2. 多 tagger 共享一遍扫描:TaggerProcessor.process_single 打开一次输入流,一行解码一次,喂给所有 tagger,再按输出路径聚合同一个文件句柄写出(python/dolma/core/runtime.py:241_make_output_streams,python/dolma/core/runtime.py:158)。IO 是这个尺度下唯一值钱的东西。
  3. 哈希分区实现「多机一个 Bloom filter」:build_hashes 先算第一个哈希,若 hash % num_partitions != partition_index 则该 key 不归本机管,直接跳过(src/deduper.rs:99)。N 台机器各跑各的,合起来等价于一台机器跑全量——不需要任何网络同步。
  4. span 替换里的 {} 占位符:replacement 字符串里 {} 会被替换成被删除的原文(src/shard.rs:402-406 附近的 replace("{}", …)),于是「删除」和「脱敏替换」共用同一套机制——'' 是删,|||EMAIL_ADDRESS||| 是脱敏。
  5. provenance 白送:mixer 给每行输出写入 原文件名:行号(src/shard.rs:455-471),训出问题样本时可以一路追回去。
  6. 随机数 tagger 即采样器:random_number_v1 给每篇文档打个 [0,1) 均匀随机数(python/dolma/taggers/sampling.py:10),之后任何比例的抽样都只是 mixer 里一条 random < p 的 exclude 表达式(实例见 configs/dolma-v1_6/sample.yaml)。

6. 边界与局限

诚实地说,这个库有明显的「AI2 内部工具」痕迹:

  • 路径约定是硬编码的:/documents//attributes/<name>/ 靠字符串替换(src/shard.rs:53src/deduper.rs:141),你的数据不照这个目录结构摆就得改代码或软链。
  • Bloom filter 有假阳性:去重判重可能误杀(unique 文档被判为重复),官方文档也明说(docs/deduplication.md);尺寸要按 estimated_doc_count 事先估好,低估会系统性抬高误杀率。
  • attributes 必须与 documents 严格逐行对齐:mixer 靠 id 校验(src/shard.rs:294-306 的 "Mismatched ids" 报错),任何一步乱序/丢行都会在混合时炸出来——这既是保护也是负担。
  • 并行模型是单机多进程,不是分布式:BaseParallelProcessormultiprocessing(python/dolma/core/parallel.py:369),Rust 侧用 threadpool;跨机靠的是「按文件分片/按哈希分区」的人工拆分,没有调度器。
  • 真实配置里混着脏数据痕迹:例如 configs/dolma-v1_6/mixing/c4.yaml 里有 $@.attributes 这种疑似笔误的 JSONPath,以及注释掉的旧规则——说明这些 YAML 是手工迭代的实验记录,不是精心维护的 API。
  • 看不出全局混洗:mixer 的 shuffle: true 实际只决定分 shard 的方式(src/mixer.rs:12-16),代码里没有跨文件的全局文档 shuffle;混洗主要发生在 tokenize 阶段(05 章)。
  • 看不出数据集本身:语料托管在 HuggingFace Hub,本仓库只有配方与工具;sources/ 目录只是来源说明文档。

7. 横向对比

同书架上与「数据」相关的兄弟项目:

  • datatrove(HuggingFace):同样是「流水线式语料清洗」,但抽象是** executor + pipeline 步骤的流式 DAG**,且原生面向 Slurm 集群调度;Dolma 的抽象是**「先标注、后混合」的两阶段**,过滤规则全部外置成数据(datatrove 的过滤器更多是代码)。Dolma 的 span 级替换(脱敏而不丢文档)是比 datatrove 的文档级过滤更细的一档。
  • hf-datasets:定位是通用数据集加载/处理库(map/filter 在内存或 arrow 上),面向「万级数据集的任意变换」;Dolma 是专用预训练语料生产线,面向「一个语料的数十亿行、一遍过」。
  • olmo:Dolma 的下游——OLMo 训练栈吃的正是 Dolma 产出的 memmap token 流(.npy + 偏移),两者通过 dolma tokens 的输出格式咬合。

一句话:datatrove 把流水线当一等公民,Dolma 把「标注缓存」当一等公民——后者让你能为一篇文档攒十几份报告再慢慢调配方,这正是做数据消融实验(ablation)最想要的形态。

8. 代码地图(导航索引)

主题文件路径符号名
并行框架入口python/dolma/core/parallel.pyBaseParallelProcessor.__call___get_all_paths_multiprocessing_run_all
断点续跑python/dolma/core/parallel.py_process_single_and_save_statusMETADATA_SUFFIX
tagger 运行时python/dolma/core/runtime.pyTaggerProcessor.process_single_write_sample_to_streams_determine_output_paths_for_taggers
tagger 基类/注册表python/dolma/core/taggers.pypython/dolma/core/registry.pyBaseTagger.tagBaseTagger.group_outputTaggerRegistry
数据模型python/dolma/core/data_types.pyInputSpecOutputSpecSpanDocResult
Bloom filtersrc/bloom_filter.rsBloomFilteroptimal_number_of_hasherssuggest_size_in_bytesinsert/contains
去重主流程src/deduper.rsrunwrite_attributesbuild_hashes
混合/装箱src/mixer.rssrc/shard.rsmixer::runShard::split_streamsShard::process
过滤表达式src/filters.rsDocFilterJqDocFilter::should_keepJsonPathFilter
span 替换src/shard.rsSpanReplacer::find_spans_to_replaceReplacement
S3/本地文件缓存src/shard.rsFileCache::prepare_inputfinalize_output
PII taggerpython/dolma/taggers/pii.pyBasePiiFilter.predict_postprocess
URL 黑名单python/dolma/taggers/url.pypython/dolma/core/url_blocker.pyBaseUrlTaggerBaseDomainTaggerUrlBlocker(Rust adblock 引擎)
质量启发式python/dolma/taggers/gopher.pypython/dolma/taggers/c4.pyGopherAttributesget_attributes
fastText 分类python/dolma/core/ft_tagger.pypython/dolma/taggers/jigsaw.pyBaseFastTextTaggerFastTextJigsawHatespeechDocumentTagger
采样python/dolma/taggers/sampling.pyRandomNumberTagger
tokenize 执行器python/dolma/tokenizer/executor.pyMemMapParallelWriter.process_singletokenize_in_parallel
memmap 写出python/dolma/tokenizer/memmap_writer.pyMemmapWriterwrite_manyclose
Python-Rust 桥src/lib.rspython/dolma/__init__.pymixer_entrypointdeduper_entrypointdolma.mixer
真实生产配方configs/dolma-v1_6/mixing/c4.yamlconfigs/dolma-v1_6/sample.yaml(YAML,非代码)