跳到主要内容

数据截至 (上游 commit 669f534823b0)

02 · 去重:Bloom filter 与三级判重

这一章讲什么: Dolma 的去重器(全 Rust):Bloom filter 怎么定尺寸、怎么落盘;文档级/段落级/ngram 级三种判重粒度;以及「按首哈希分区」这个让多机去重不需要网络同步的技巧。


1. 它要解决的小问题

网页语料里重复泛滥:同一篇文章被转载几百次、同一个模板页生成几百万个 URL、大段文字被复制粘贴。

精确去重需要记住所有见过的内容——几十亿篇文档,哈希集合的内存开销谁也扛不住。而且去重是全局状态操作:机器 A 扫过的文档,机器 B 也要知道。

Dolma 的选择:用 Bloom filter 拿「少量误杀」换「内存可行 + 并发可行」


2. 思路:三个先验判断

为什么用 Bloom filter? 它是「可能说有的有、绝不说没有的没有」的概率型集合:查询返回 false 一定没见过,返回 true 可能误判(假阳性)。去重场景里假阳性 = 误杀一篇无辜文档,代价可接受;换来的是内存比哈希集合小一个数量级以上,还自带可落盘的位图形态。

为什么并发安全? 判重逻辑是「contains 为 false 才 insert」。两个线程同时查到 false、同时插入,结果只是都写一遍相同的位——幂等,不需要锁。底层位数组用 AtomicU32fetch_or,连写冲突都没有。

为什么能多机并行? 按 key 的第一个哈希值取模,把 key 空间静态切成 N 份,每台机器只处理自己那份 key。相同 key 必然落到同一台机器,所以分区间不需要任何通信(细节见 §5.4)。

三种判重粒度

粒度配置key 是什么产出
文档级dedupe.documentsJSONPath 取出的字段(典型是 $.metadata.url)整文档 span [0, len, 1]
段落级dedupe.paragraphs整段文本(默认按 \n 切)每个重复段落一个 span
ngram 级dedupe.paragraphs.by_ngram段内每个 n 词滑窗重叠率 ≥ 阈值的段落 span,score 是重叠率

三者互斥的是前两者(入口处 XOR 校验,src/deduper.rs:29-32);ngram 是段落级的子模式。


3. 图示:一次判重的流向

documents/x.jsonl.gz 逐行


取 key(URL / 段落 / ngram)


首哈希 % N == 本机分区号? ──否──► 跳过(不归我管)
│是

contains(全部 k 个哈希)? ──是──► 标 span(重复)
│否

insert(把 k 个位置 1) → 放行(首次出现)

怎么读: 每个 key 先过「分区闸」再过「判重闸」。read_only 模式下 insert 是空操作——整个 filter 变成一张只读黑名单,这正是去污(decontamination):先拿评测集建成 filter,再以只读模式扫训练语料,命中即标泄漏。


4. 原理演示

# 示意,非源码
class BloomFilter:
def __init__(self, num_bits, seeds):
self.bits = bytearray(num_bits // 8)
self.seeds = seeds # 每个种子派生一个独立哈希函数

def _hashes(self, key):
return [hash_with_seed(key, s) % len(self.bits) for s in self.seeds]

def insert(self, key):
for h in self._hashes(key):
self.bits[h] = 1

def contains(self, key):
return all(self.bits[h] for h in self._hashes(key))

# 判重主循环:见过 → 标重复;没见过 → 记住并放行
for doc in docs:
key = extract_key(doc) # 如 doc.metadata.url
if bf.contains(key):
mark_duplicate(doc, span=(0, len(doc.text), 1.0))
else:
bf.insert(key)

重点看 containsinsert 用的是同一组哈希——k 个位全 1 才算见过,所以单个哈希冲突不会误判,假阳性率随 k 和位数组大小指数下降。

尺寸怎么定?教科书公式:给定预期元素数 n 和目标假阳率 p,最优哈希数 k = (m/n)·ln2(m 为位数)。Dolma 的做法是翻倍搜索:从 1MB 起,尺寸不断 ×2,直到代入公式的假阳率 ≤ 目标(suggest_size_in_bytes,src/bloom_filter.rs:45-58);k 由 optimal_number_of_hashers 现算(src/bloom_filter.rs:27-32)。用户只需给 estimated_doc_count + desired_false_positive_rate 两个业务数字。


5. 真实实现

5.1 位数组与哈希族

BloomFilter(src/bloom_filter.rs:15)的三件东西:

  • bits: Vec<AtomicU32>——位数组按 32 位字打包。定位某一位:index = hash / 32 % len,bit = hash % 32(src/bloom_filter.rs:218-220src/bloom_filter.rs:229-231 的同一公式,insert/contains 对称)。
  • hash_builders: Vec<RandomState>——k 个 ahash 哈希器,每个带 4 个 u64 随机种子(src/bloom_filter.rs:78-86)。种子必须随 filter 一起持久化,否则读回的位图无法复算哈希。
  • read_only: bool——insert 开头直接判空操作(src/bloom_filter.rs:215-216)。

写入用 fetch_or(1 << bit, Ordering::Relaxed)(src/bloom_filter.rs:221):原子或,无锁、幂等。contains 任一位为 0 即返 false(src/bloom_filter.rs:225-234)。

5.2 落盘格式

write_to_file(src/bloom_filter.rs:151)写一个简单二进制:魔数 0x81F0F117 + 版本号 + 哈希器种子序列 + 位图原始字节(经 unsafefrom_raw_parts 直接转字节切片,src/bloom_filter.rs:170-174)。from_file(src/bloom_filter.rs:100)校验魔数与版本后原样读回。整个去重状态就是一个文件——这就是「先去污建库、再只读扫语料」这类玩法的物理基础。

5.3 段落级与 ngram 级判重

write_attributes(src/deduper.rs:116)是单文件处理主体:

  • 文档级:JSONPath 取 key(src/deduper.rs:214-227jsonpath_rust::JsonPathFinder),contains 命中则写 [0, len, 1] 整文档 span(src/deduper.rs:294-301),否则 insert(src/deduper.rs:303-305)。
  • 段落级:按分隔符切段落(src/deduper.rs:322,默认 \n),手工维护字符偏移得到 [par_start, par_end],逐段查 filter(src/deduper.rs:368-378)。
  • ngram 级:段落再切成 unicode 词边界(src/wimbd/tokens.rs:11tokenize,基于 unicode_segmentationsplit_word_bounds),滑动窗口维护一个 VecDeque 当 ngram(src/deduper.rs:385-416);统计「命中 ngram 数 / 总 ngram 数」,≥ overlap_threshold 才把整段标为重复,score 就是重叠率(src/deduper.rs:454-465)。stride 控制跳采样,避免每个位置都算哈希。
  • 短段落回退:ngram 数不足 2 时退化为整段判重(src/deduper.rs:421-449),除非配置 skip_short_paragraphs

5.4 哈希分区:多机去重零协调

两级分区,都在 src/deduper.rs:

  • 按文件分区(粗):file_partition 开启时,文件路径哈希取模不属于本机就直接 continue(src/deduper.rs:44-49)——把文件清单静态切开。
  • 按 key 分区(细):build_hashes(src/deduper.rs:99-113)只先算第一个哈希,hash % num_partitions != partition_index 就返回空——这个 key 的判重权归别的机器。命中分区才补算剩余 k-1 个哈希,省 CPU。

相同 key 的所有副本必然哈希相同、落到同一台机器,所以 N 台机器各扫全集的一部分文件,合并效果等价于单机全扫。这就是 Dolma 在 README 里宣称「billions of documents concurrently」的去重侧答案:没有分布式框架,只有取模。

5.5 产出仍是「报告」

deduper 不删数据,只写 attributes(路径同样是 /documents//attributes/<name>/ 替换,src/deduper.rs:142-144,且替换失败会直接 panic 防误写原文目录,src/deduper.rs:147-153)。段落级产出示例:

{"id": "...", "attributes": {"bff_duplicate_paragraph_spans": [[120, 980, 1.0]]}}

删不删、怎么删,交给第 03 章的 mixer 按 score 阈值决定。若开启分区,属性名还会带 _<分区号> 后缀(src/deduper.rs:130-135),由后续合并配置分别引用。


6. 关键细节与坑

  • 假阳性 = 误杀,且随数据量上涨:filter 尺寸按 estimated_doc_count 一次定死,实际数据量显著超过估计时,假阳率会高于目标值——估大别估小。
  • 分区模式下属性名带后缀(bff_duplicate_docs_0 之类):mixer 的过滤表达式要按后缀名逐个写,真实配置里能看到这种一长串。
  • read_only 文件必须已存在:Python CLI 侧若远端 bloom filter 下载不到且是只读模式,直接抛错(python/dolma/cli/deduper.py:220-228)——否则等于拿空 filter 去污,什么都查不出来还自以为查过。
  • offset 手工记账:段落 span 的偏移是字符数累加 + 分隔符补 1(src/deduper.rs:324-331),换自定义 paragraph_separator 时这段逻辑仍按 1 个字符补——多字符分隔符下 span 会错位(inferred,代码未见对分隔符长度的处理)。
  • 每篇文档输出一行 attributes,哪怕为空数组:保证与 documents 逐行对齐(第 01 章的约定),缺行会让 mixer 的 id 校验炸掉。
  • 看不出跨分区的一致性保证:若同一 key 的副本所在文件被「按文件分区」切到不同机器,而 key 分区又没开,理论上会漏判;两种分区的组合用法代码里没有防呆(inferred)。