跳到主要内容

数据截至 (上游 commit a649de79c14a)

04 · 去重

这一章讲什么: 全库工程含量最高的部分。以 MinHash 四阶段为主线(FineWeb 生产配置),讲清「签名 → 桶内归并 → 并查集聚类 → 过滤」这条流水线如何在几十亿文档上做近似去重;然后看三个变体:句子级精确去重(会改文本)、整文精确去重、Bloom filter。


1. 它要解决的小问题

网页语料里同一篇文章会被转载几百次。训练时见过几百次的东西会被模型背下来,所以要去重。难点分层:

  • 近似:转载时有小编辑(改日期、加导航),精确哈希抓不到——要 MinHash 这类相似度指纹;
  • 规模:几十亿文档,两两比对是 O(n²),内存也装不下全量指纹;
  • 粒度:有时整篇重复,有时只有中间几句重复(C4 的处理是「删掉重复的句子」而非弃文)。

2. 思路:把「全局比对」拆成「落盘的排序 + 归并」

先看全局直觉。MinHash 的经典 trick 是:把文档切成词的 n-gram(shingle),每个 shingle 算哈希,取最小的若干个作为签名——两篇文档签名撞得越多,Jaccard 相似度越高的概率越大。再把 112 个签名值切成 14 桶 × 8 个,任一桶内 8 个值全同就报疑似重复。

DataTrove 的工程贡献是:不发明新的分布式比对框架,而是把问题规约成两件古老的事——排序归并排序

  1. 每个 rank 给本地文档算签名,按签名排序后落盘(二进制定长记录);
  2. 同一桶内,多个已排序文件做 k 路堆归并——相同签名在归并流里必然相邻,扫一遍就找到所有重复对;
  3. 重复对做并查集聚类,每个簇只留代表元,其余写进 .remove 清单;
  4. 过滤阶段按清单逐文档丢弃。

全程顺序 I/O,内存占用可控,中间产物(签名、重复对、清单)全是可复算的落盘文件。


3. 归一化与签名:让「小编辑」失效

3.1 文本先归一化

比对前先抹掉无关差异。simplify_textsrc/datatrove/utils/text.py:212-257)按 TextNormConfig:185-193)依次做:小写化 → 数字统一替换成 0NUMBERS_PATTERN\p{Nd} 匹配任意文字的数字,:202-205)→ 标点换成空格 → 空白压缩 → NFD 拆变音符并去掉。所以「2023 年 1 月」和「2024 年 3 月」转载版会算出同一批 shingle。

3.2 签名的数学

MinhashDedupSignaturesrc/datatrove/pipeline/dedup/minhash.py:124)。默认配置(MinhashConfig:40-57):

参数默认含义
n_grams55 个词一个 shingle
num_buckets14桶数(LSH 的 band 数)
hashes_per_bucket8每桶哈希数
hash_config64-bit xxhashshingle 的基础哈希;FineWeb 换 sha1 求稳(examples/fineweb.py:80-88

文件头注释直接给了概率推导(:30-35):14 桶 × 8 哈希对应真实相似度阈值 (1/14)^(1/8) ≈ 0.72,相似度 0.8 的文档被报出的概率约 0.924。

签名的计算是教科书式的 universal hashing(get_signature:172-188):预生成随机参数 (a, b)parameters:154-170,种子固定保证全集群一致),phv = (shingles * a + b) % (2^61 - 1),取每个哈希函数的最小值,再切成 14 组。梅森素数 2^61 - 1 定义在 :28

3.3 落盘 + 排序 + 校验

每个 rank 为每个桶写一个 bucket_{bi:03d}/{rank:05d}.minhash.sig 文件,记录是定长二进制:8 个哈希 + 文档序号(struct.pack:255-261)。写完后读回、用 numpy 结构化 dtype 按全字段排序、写回:266-297)。

这里有一段不常见的防御(:290-321):排序后的字节先算 blake2b 摘要,写盘后再读回重算比对,不一致就抛「签名损坏」。注释明说动机——「未排序的签名是 MinHash 的头号事故源」:292)。因为下一阶段(堆归并)的正确性完全建立在「每个输入文件已排序」上,一旦有个文件悄悄没排好,结果是静默的错误去重,比 crash 可怕得多。


4. stage 2:桶内 k 路堆归并

MinhashDedupBucketsminhash.py:324)负责找重复对。它的并行方式值得细看——两级切分

  • 按桶切:rank 除以每桶 worker 数得到桶号(divmod(rank, workers_per_bucket):395);要求 world_size % num_buckets == 0:393);
  • 桶内再按哈希值区间切:get_worker_hash_range:360-389)取第一个签名文件,按行数均分,读出边界行的哈希值作为本 worker 的 [hash_min, hash_max)。因为文件已排序,行均分 ≈ 值域均分。

于是每个 worker 只对「桶 bi 里哈希落在某区间的签名」负责。worker 打开该桶的所有签名文件,各配一个 read_sigs 生成器(:84-121,内部用 seek_to_start 二分定位到 hash_minsrc/datatrove/utils/binaryio.py:54-94),然后用 heapq 做 k 路归并(:454-495):

# src/datatrove/pipeline/dedup/minhash.py:467-475(精简)
while pq:
v = heapq.heappop(pq)
if not v.is_from_index():
if last is not None and last.sig == v.sig:
# 写 (file_id1, doc_id1, file_id2, doc_id2) 四元组
out_f.write(struct.pack("<4I", ...))

归并流里相同签名必然相邻,所以 last.sig == v.sig 就是全部判定逻辑——排序把「全局比对」变成了「相邻比较」。

index 机制:去重还可以「对着一个既有数据集做」——先用 MinhashBuildIndex:691)把参考数据集的签名合成 index 文件,stage 2 把 index 文件也并入归并流(:426-452),命中 index 的重复对用 (SENTINEL, SENTINEL):37)占位写出(:473-476),下游据此知道「这篇撞了参考集」。


5. stage 3:并查集聚类

MinhashDedupClusterminhash.py:500)把重复对聚成簇:A-B 重复、B-C 重复 ⟹ A、B、C 一个簇,只留一个。

实现是教科书并查集 + 两个优化(parent/union:537-558):路径压缩(union_set[x] = parent(union_set[x]))和按大小合并(小的挂到大的下,保持树浅)。SENTINEL 簇特殊照顾——撞 index 的文档全挂到 SENTINEL 根下(:553)。

然后一次性遍历(:574-596):非根节点写进 {file:06d}.remove(待删),根节点是幸存代表;可选写 .clusters(簇 id)和 .sizes(簇大小)供分析。注意它 必须单任务跑assert world_size == 1:533)——重复对全集要装进内存,这是整条链唯一的单点,FineWeb 给它 8 核 200GB(examples/fineweb.py:150-156)。

6. stage 4:按清单过滤

MinhashDedupFilterminhash.py:599)回到文档流上。它打开本 rank 的 {rank:06d}.remove,文件里是排序好的待删 doc_idx,于是双指针扫过文档流:idx 命中就丢(或写进 exclusion_writer),否则放行(:641-688)。还能顺手把 minhash_cluster_id / minhash_cluster_size 写进 metadata(:653-677)——下游可以按簇大小做数据加权。

整条四阶段链在 FineWeb 里的组装方式见 examples/fineweb.py:103-180:四个独立 SlurmPipelineExecutor,用 depends= 串成 DAG,tasks 分别是 1000(签名)/ 700 = 14 桶 × 50(归并)/ 1(聚类)/ 1000(过滤)。


7. 变体一:句子级精确去重(会改文本)

C4 论文的做法(文件头注释原文引用,:1-4):「数据集中出现超过一次的连续三句,只保留一份」。这是删不是删

src/datatrove/pipeline/dedup/sentence_dedup.py 三阶段:

  1. SentenceDedupSignature:68):每篇文档切句 → 连续 3 句拼成一个 span → 归一化后哈希,签名是 (hash, doc_id, sent_idx) 三元组(get_hashes:133-146)。落盘时按哈希值区间分桶给 finder(save_hashes:100-131),所以 stage 2 可以多任务并行;
  2. SentenceFindDedups:189):同一套「排序 + 堆归并 + 相邻比较」,输出 (doc_id, sent_idx) 对;
  3. SentenceDedupFilter:294):remove_dup_sentences:329-378)用 span_tokenize 的字符区间定位句子,命中的 span 起点把 drop_until 设为 idx + n_sentences——重叠的重复 span 会连成一片删掉;删完若剩余词数/句数低于下限(默认 50 词 / 3 句,SentDedupConfig:41-50)则弃文。exclusion_writer 存的不是原文而是带 >>>/<<< 标记的「删改对照版」(:352-368),方便人工检查删了什么。

8. 变体二与三:整文精确去重、Bloom filter、ExactSubstr

  • Exact dedupsrc/datatrove/pipeline/dedup/exact_dedup.py):整文哈希版,同样三阶段。独特点是 priorityExactDedupConfig.document_priority:29-42)给文档打 1~65535 的优先级,签名排序按 (hash, -priority, doc_id)save_hashes:102-115)——同一哈希簇里优先级最高的排最前、被保留,其余进删除清单。适合做「同文多源,留权威源」。stage 3 还会给幸存文档写 duplicate_count metadata(:415)。
  • Bloom filtersrc/datatrove/pipeline/dedup/bloom_filter.py:66):单 block、无多阶段——每个文档的 shingle 哈希进布隆过滤器,按「已见过的 shingle 占比 ≥ duplicate_threshold(默认 0.8)」判重。构造时还会用公式估算假阳性率并警告(:102-107)。省内存、一次过,代价是有假阳性且不支持 index。
  • ExactSubstrsrc/datatrove/pipeline/dedup/exact_substrings.py):子串级去重,包装 google 的 deduplicate-text-datasets(suffix array),类是 ESDatasetToSequence:42)/ESMergeSequences:85)/ESRangeRemover:149)。依赖外部 Rust 工具链,本文不展开。

9. 关键细节 / 坑

  • 签名排序是性命攸关的不变量。 stage 2 的归并对每个输入断言单调性(minhash.py:111-113:469:494),乱序直接抛错并建议重跑该文件。stage 1 的 blake2b 自检(:290-321)就是为这道不变量上的保险。
  • finder_workers 与任务数必须对齐。 sentence/exact dedup 的 stage 1 按哈希区间给 stage 2 预分桶,stage 2 的 tasks 数不等于它就直接报错(sentence_dedup.py:222-228)。
  • doc_idx 只是 rank 内序号。 所有清单都按 (rank 对应的文件, doc_idx) 定位(minhash.py:580{file:06d}.remove),所以 stage 4 必须用和 stage 1 完全相同的 reader 配置重读数据——FineWeb 里就是同一个 INPUT_READER 对象复用(examples/fineweb.py:97-100:163)。输入变了,清单就对错文档。
  • 去重前的过滤器别乱序。 MinHash 的归一化(simplify_text)只抹格式差异;若上游过滤器会改文本(C4 抠行),签名阶段读到的就是改后的文本——这是符合预期的,但要意识到顺序即语义。
  • 代码里看不出的:14 桶 × 8 哈希这套参数在超大规模下的实际假阳/假阴率、以及 hashes_per_bucket 改小一档对召回的影响,只能查 FineWeb 技术报告的消融,代码里只有 :30-35 的概率估算注释。

10. 本章代码地图

主题文件路径符号名
MinHash 配置src/datatrove/pipeline/dedup/minhash.pyMinhashConfigHashSigSENTINEL
签名(stage 1)同上MinhashDedupSignature.get_signatureget_shinglescheck_can_skip_sig_writing
桶内归并(stage 2)同上MinhashDedupBuckets.runget_worker_hash_rangeread_sigs
并查集聚类(stage 3)同上MinhashDedupCluster.run(内嵌 parent/union
按清单过滤(stage 4)同上MinhashDedupFilter.run
建 index同上MinhashBuildIndex
句子级去重src/datatrove/pipeline/dedup/sentence_dedup.pySentenceDedupSignature.get_hashesSentenceDedupFilter.remove_dup_sentences
整文精确去重src/datatrove/pipeline/dedup/exact_dedup.pyExactDedupSignatureExactFindDedupsExactDedupFilter
Bloom filtersrc/datatrove/pipeline/dedup/bloom_filter.pySingleBloomFilterget_optimal_k
归一化src/datatrove/utils/text.pysimplify_textngrams
哈希函数src/datatrove/utils/hashing.pyHashConfigcreate_hash_func
二分定位src/datatrove/utils/binaryio.pyseek_to_startread_tuples_from_file