数据截至 (上游 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-55 的 input.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-143的Selector::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)。