数据截至 (上游 commit 669f534823b0)
05 · tokenize 管道:环形混洗与 memmap 落盘
这一章讲什么:
dolma tokens怎么把混合好的 JSONL 变成训练用的 token 流:为什么「ring buffer + 局部 shuffle」能代替全局混洗、memmap 文件怎么边写边建、元数据 CSV 记什么。这是 Dolma 流水线的最后一环,出口正对 OLMo 训练栈。
1. 它要解决的小问题
混合好的语料是几千个 JSONL 文件、几十亿篇文档。进训练前还差两步:
- tokenize:逐篇编码成 token id,拼成连续流,文档之间插 EOS 分隔。
- 混洗:训练框架按顺序读 token 流,如果同一来源的文件连在一起,模型会在一个 batch 里看到整批同质数据,训出来偏。
全局 shuffle 几十亿篇文档需要外部排序,又贵又麻烦。Dolma 的问题因此是:能不能用「近似混洗」达到「足够均匀」?
2. 思路:三个便宜的随机化叠加
Dolma 的答案是叠加三层廉价的随机性,凑出接近全局混洗的效果:
| 层 | 动作 | 成本 |
|---|---|---|
| 文件顺序 | 启动时把全部输入路径 shuffle 一次 | O(文件数) |
| 跨文件轮转 | 每个进程同时开 8 个文件的 token 流,轮流(或按大小加权)各取一篇 | 常数内存 |
| 局部洗牌 | 攒满 1 万篇后 random.shuffle,再落盘 | O(1 万) |
直觉:只要来源在盘上分布得足够开,训练时的「顺序读」就近似「随机读」。环形缓冲让任意时刻的读取位置分散在 8 个不同文件 上;局部洗牌消灭窗口内的残余顺序;文件级 shuffle 保证窗口与窗口之间不相关。
docs/tokenize.md 也明说这叫 "light shuffling",并提醒:想混得更匀,优先加大 ring 里的文件数,而不是加大局部洗牌窗口(python/dolma/tokenizer/executor.py:141-146 的注释同理)。
3. 图示:一个进程的 token 流
文件 A ─┐
文件 B ─┤ ring(8 个生成器)
文件 C ─┼─► next(ring[j]) 轮流取一篇 ─► accumulator(攒 1 万篇)
┘ ▲ │
读完即踢出,补新文件 random.shuffle
▼
MemmapWriter 顺序写
part-00000-00000.npy (+ .csv.gz)
写满 500M token → 开下一卷
怎么读: 左边 8 个文件生成器是「进水口」,中间 1 万篇的蓄水池是「洗牌间」,右边 memmap 是「出水口」。N 个进程各自独立跑同一套,各自写自己的 part 文件。
4. 原理演示
# 示意,非源码
def tokenize_bucket(paths, writer, ring_size=8, local_shuffle=10_000):
ring = [tokenize_file(paths.pop()) for _ in range(ring_size)] # 每个元素是逐篇产 token 的生成器
accumulator = []
while ring:
for i in range(local_shuffle):
j = i % len(ring) # 轮转;也可按文件大小加权抽样
try:
accumulator.append(next(ring[j]))
except StopIteration: # 这卷读完了
ring.pop(j)
if paths: ring.append(tokenize_file(paths.pop()))
random.shuffle(accumulator) # 窗口内洗牌
for seq in accumulator:
if not writer.write(seq): # 当前 memmap 写满
writer = writer.next_volume() # 换卷继续
accumulator = []
重点看 j = i % len(ring):不是随机抽,是确定性地轮流转——轮转保证每个文件被均匀消费,随机性由「文件顺序已 shuffle + 窗口内洗牌」提供。
5. 真实实现
5.1 分桶:先打散文件,再按进程数切
MemMapParallelWriter.__call__(python/dolma/tokenizer/executor.py:252)覆盖基类行为:
- 全部输入路径先
random.shuffle(python/dolma/tokenizer/executor.py:257-258)。 - 按
step_size = 总文件数 / 进程数切成桶,每桶是一个进程的输入(python/dolma/tokenizer/executor.py:262-274);并做三条一致性校验(无空桶、桶数 ≥ 进程数、文件总数守恒,python/dolma/tokenizer/executor.py:278-289)。 - 传给子进程的不是路径而是桶下标(
python/dolma/tokenizer/executor.py:292-293的source_indices)——路径清单太长,pickle 进进程间通信不划算,子进程用下标自取。
5.2 环形缓冲:轮转与按大小加权
process_single(python/dolma/tokenizer/executor.py:46)的读取循环:
- 初始化 ring 装
min(ring_size, len(paths))个tokenize_file生成器,默认ring_size = 8(python/dolma/tokenizer/executor.py:54,116-126)。 - 每个窗口攒
local_shuffle(默认 10_000,python/dolma/tokenizer/executor.py:53)篇;sample_ring_prop=False时j = i % len(tokenizer_ring)轮转(python/dolma/tokenizer/executor.py:173),为 True 时按文件大小加权np.random.choice(python/dolma/tokenizer/executor.py:169-171,概率随 ring 成员变动即时重算,python/dolma/tokenizer/executor.py:208-209)。 - 某个文件读完(
StopIteration)立即踢出 ring、补新文件(python/dolma/tokenizer/executor.py:185-206)。 - 攒满即
random.shuffle(accumulator)(python/dolma/tokenizer/executor.py:219)后交 writer;写 不下的余量重洗后写进下一卷(python/dolma/tokenizer/executor.py:223-244)。
进度上报同样有「队列积压就翻倍间隔」的自适应(与第 01 章 tagger 相同的手法,python/dolma/tokenizer/executor.py:211-216)。
5.3 memmap:定长预分配,关卷裁尾
MemmapWriter(python/dolma/tokenizer/memmap_writer.py:19)的写法:
- 打开时按
max_tokens(默认 500M token,约 1GB,DEFAULT_MAX_TOKENS,python/dolma/tokenizer/memmap_writer.py:22)用np.memmap(mode="w+")预分配(python/dolma/tokenizer/memmap_writer.py:153-156)。 write把一篇文档的 token 顺序拷进当前偏移,同时往 CSV 写一行元数据(python/dolma/tokenizer/memmap_writer.py:82-92):id、src(源文件)、loc(源文件行号)、start、end——token 流里的每一段都能找回原文位置,这是数据可追溯的又一层。- 关卷时若没写满,重开一个精确大小的 memmap 把有效部分拷过去再删旧文件(
python/dolma/tokenizer/memmap_writer.py:179-188)——先占坑再裁尾,避免边写边扩容的反复拷贝。 - dtype 默认 uint16,也可由
Tokenizer按词表大小自动选(np.min_scalar_type(vocab_size - 1),python/dolma/tokenizer/tokenizer.py:105)。 - 目标是 S3 时先写本地临时文件,关闭时上传(
python/dolma/tokenizer/memmap_writer.py:191-201)。