跳到主要内容

数据截至 (上游 commit 669f534823b0)

03 · 混合器:一遍 IO 完成合并、过滤与文本手术

这一章讲什么: dolma mix 的全流程——装箱、按行归并属性、表达式过滤、span 替换的状态机、provenance 落盘。这是整条流水线里唯一写「最终数据」的环节,前两章攒的所有「报告」都在这一章兑现。


1. 它要解决的小问题

到混合这一步,你手里有:原文 documents,以及十几份 attributes(gopher 信号、去重 span、PII 位置、随机数……)。

你要的是:按一张 YAML 配方,把这些合成最终训练文件。具体是三件事:

  • 合并:把同一篇文档的所有属性拼到它身上(或直接用于判决)。
  • 判决:整篇的留/弃(过滤),以及篇内某些段的删/换(span replacement)。
  • 打包:输出成大小均匀的 shard 文件,并且知道每行来自哪。

难点在规模:这些事要在几十亿行上做,且只能扫一遍


2. 思路:判决是表达式,手术是状态机

两个关键设计:

过滤不写代码,写表达式。 一条 exclude 规则就是一个 JSONPath(或 jq)表达式,对着「文档 + 合并后属性」的 JSON 求值,命中即弃。于是「调配方」= 改 YAML,不需要重新跑任何 tagger。

span 替换不重解析,走字符流。 要删/换的区间先按 start 排序,然后一次字符遍历,边读边决定「这个字符抄不抄进新文本」——是个小型状态机。

配方长什么样(真实节选)

configs/dolma-v1_6/mixing/c4.yaml 的结构(删减):

streams:
- name: c4
documents: [s3://…/c4/v0/documents/train/*.gz]
attributes: [olmo_mix_v1_taggers, dedupe_paragraphs, dedupe_docs, tokenizer_repetitions_v2r2]
filter:
exclude:
# 词数 < 50 的丢弃(查 gopher 报告里的测量值)
- "$.attributes[?(@.gopher_rules__gopher_v1__word_count[0][2] < 50)]"
# 整篇重复的丢弃(查去重报告)
- "$.attributes[?(@.bff_duplicate_docs[0][2] >= 1.0)]"
span_replacement:
# 重复段落:score ≥ 0.5 的删掉
- span: "$.attributes.bff_duplicate_paragraph_spans"
min_score: 0.5
replacement: ''
# 邮箱:score ≥ 0.5 的换成占位符
- span: "$.attributes.pii_detection__pii_regex_with_counts_fast_v2__EMAIL_ADDRESS"
min_score: 0.5
replacement: " |||EMAIL_ADDRESS||| "

读法:exclude 是「整篇枪毙」,span_replacement 是「篇内手术」;两者都只是对第 01/02 章攒下的 [start, end, score] 测量值做阈值查询。


3. 图示:一行数据的混合之旅

documents 行 attributes 行 ×N(同 id)
│ │
└──────► 按行对齐归并(id 校验)◄──────┘


DocFilter.should_keep?
(include/exclude 表达式)
│ false → 丢
▼ true
SpanReplacer ×M
按分数筛出待替换 span


逐字符重写 text(删/换)


丢字段 + 长度复查 + provenance


写入输出 shard

怎么读: 每个输入文件经过一次「归并 → 判决 → 手术 → 盖章」;所有 attr 文件与 doc 文件是同步推进的迭代器,任何一份行数对不上都会触发 id 校验错误。


4. 原理演示

# 示意,非源码
def mix_one_line(doc_line, attr_lines, filter_expr, replacers):
doc = json.loads(doc_line)
# ① 归并:每份 attributes 的 key 直接并入(撞名后到的覆盖)
for line in attr_lines:
attr = json.loads(line)
assert attr["id"] == doc["id"] # 防错位
doc.setdefault("attributes", {}).update(attr["attributes"])

# ② 判决:表达式命中 → 整篇丢
if any(eval_jsonpath(e, doc) for e in filter_expr.exclude):
return None

# ③ 手术:收集待替换 span,按分数过滤,逐字符重写
spans = sorted(
(s for r in replacers for s in r.find(doc) if r.min_score <= s.score < r.max_score),
key=lambda s: s.start,
)
doc["text"] = rewrite(doc["text"], spans) # span 内字符不抄,边界处插入 replacement

# ④ 盖章
doc["metadata"]["provenance"] = f"{filename}:{line_no}"
return doc

5. 真实实现

5.1 装箱:按字节数切 shard

Shard::split_streams(src/shard.rs:40)把每个 stream 的输入文件按声明的 max_size_in_bytes 贪心装箱:逐个累加文件大小,超限就封箱开新 shard(src/shard.rs:69-96)。输出名形如 {output.path}/{stream名}-{序号:04}.jsonl.gz

每个输入文件的 attributes 路径在装箱时就靠字符串替换算好(src/shard.rs:52-55input.replace("/documents/", &attr_prefix))——第 01 章目录约定的消费点。注释明说大小是近似的(不含 attributes 体积、也不算过滤后缩水,src/shard.rs:36-39)。shuffle: false 时走 split_streams_unshuffled(src/shard.rs:132):一文件一 shard、保持原目录结构,用于「只过滤不重排」的场景。

5.2 归并:同步迭代器 + id 校验

Shard::process(src/shard.rs:188)对每个输入文件:

  • 为 doc 和每份 attr 各开一个 reader(src/shard.rs:230-253),逐行同步推进。
  • 每行先比 id:attr_data["id"] != data["id"] 直接报错返回(src/shard.rs:294-306,"Mismatched ids"),属性不是对象也报错(src/shard.rs:308-316)。
  • attr 读失败或提前结束时只告警、不中断(src/shard.rs:317-340),该行的这部分属性按缺失处理——判决表达式查不到对应 key 自然不命中(inferred:这是设计上的容错取向)。
  • 全部 attr 的 key 并入 data["attributes"](src/shard.rs:342-351)。

5.3 判决:两种语法的过滤器

DocFilter 是个三选一枚举(src/filters.rs:452-458):jq / JSONPath / 无过滤(AllowAll)。

  • JSONPath 版(JsonPathFilter.should_keep,src/filters.rs:395):include 任一命中即留,exclude 任一命中即弃;include 为空默认全留。
  • jq 版(JqDocFilter.should_keep,src/filters.rs:327):表达式先编译成 jaq 的 Filter,结果按类型取真值——bool 直取、null 为假、数/字符串/数组/对象按「非空为真」(evaluate_match,src/filters.rs:299-315)。jq 能写更复杂的逻辑(比如 .attributes.foo | add),JSONPath 只够做存在性/比较查询。

调用点在 src/shard.rs:367-370:should_write 为 false 的行直接不进输出。

5.4 手术:span 替换的逐字符状态机

留下的行进入替换(src/shard.rs:371-425):

  • SpanReplacer::find_spans_to_replace(src/shard.rs:647-674)按配置的选择器(jq 或 JSONPath,src/filters.rs:131-143Selector::new)取出 span 数组,逐个检查 min_score ≤ score < max_score
  • replacement 字符串若以 $ 开头则当 jq 选择器从文档里取值,否则当字面量(Replacement::new,src/shard.rs:613-624)。
  • 重写主循环(src/shard.rs:379-425)把 span 按 start 排序后,用 char_indices() 遍历原文:span 内的字符不抄;走到 span 结束位置时插入 replacement;{} 会被替换成被删掉的原文(src/shard.rs:402-406)。所以 replacement: '' 是删除,' |||EMAIL_ADDRESS||| ' 是脱敏,'{}"' 这种还能给原文加引号——同一套机制三种用途。

逐字符而不是按字节切片,是因为 span 的 start/end 是字符偏移(Python tagger 按 Python 字符串语义产出),而 Rust 字符串按字节索引——char_indices 同时拿到两种坐标做换算。

5.5 收尾:丢字段、长度复查、provenance

  • discard_fields 指定的顶层字段被移除(src/shard.rs:427-429)——生产配置里通常把庞大的 attributes 丢掉,只留正文。
  • 替换后文本 trim 长度再查一次 min_text_length(src/shard.rs:432-434):手术把全文切光的文档仍会被丢。
  • 写入前盖 provenance:"原文件名:行号" 塞进 metadata.provenance(src/shard.rs:455-471)。
  • 落盘走 FileCache(src/shard.rs:713):S3 输出先写本地临时文件再上传,成功后留一个空文件当「此 shard 已完成」的标记(finalize_output,src/shard.rs:793-820)——与第 01 章 .done.txt 同宗的断点思路。

5.6 驱动的两端

  • Python 侧 MixerCli.run(python/dolma/cli/mixer.py:88)做配置校验(必须给 filter 或 span_replacement 之一;span_replacement 必须给 min/max_score 之一;路径必须能 glob 到文件),然后把 dict 配置扔给 Rust 的 mixer_entrypoint(src/lib.rs:35)。
  • Rust 侧 mixer::run(src/mixer.rs:11)一个 threadpool 并发处理所有 shard;输出文件已存在就跳过(src/mixer.rs:19-23)——shard 级断点续跑。

6. 关键细节与坑

  • attributes 行号必须严格对齐:归并是纯位置性的,id 校验只是事后保险。上游任何一步(自定义 tagger、手工处理)打乱了行序,mixer 会在中途炸掉,而且之前写出的 shard 不会自动清理(重跑靠「已存在即跳过」,所以半成品要手动删)。
  • jq 与 JSONPath 语义不等价:同一表达式两种语法可能一个报错一个静默不命中;真实配置里还有 $@.attributes[...] 这种疑似笔误(configs/dolma-v1_6/mixing/c4.yaml),错了不炸,只是规则不生效——配方需要小样本人工验证。
  • exclude 是「任一命中即弃」:多条 exclude 是 OR 关系,且求值在属性合并之后——顺序无关,但规则之间的交互(比如两条各自合理的阈值叠加后误杀率翻倍)只能自己算。
  • span 重叠时行为微妙:替换循环假定 span 按 start 有序、基本不重叠;重叠 span 会走 while replacements[span_index].start < i 的追跳逻辑(src/shard.rs:410-414),合并语义不直观——给同一文本配多个 replacer 时要小心区间交叉。
  • 装箱大小是名义值:输出按输入文件字节数估,过滤后实际文件会比 max_size_in_bytes 小,且各 shard 大小不一。
  • 看不出属性 key 冲突的处理:多份 attributes 有同名 key 时后者覆盖前者(src/shard.rs:343-350 的 insert 语义),配方里 attribute 列表的顺序因此是有意义的(inferred)。