跳到主要内容

数据截至 (上游 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-293source_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=Falsej = 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)。

5.4 tokenizer 适配层

Tokenizer(python/dolma/tokenizer/tokenizer.py:72)是 HF tokenizers 库的薄封装:支持 fast/slow 两种后端、截断方向、是否在段落级先切分(segment_before_tokenization,针对某些 tokenizer 长文档极慢的问题)、以及必须至少给 bos_token_ideos_token_id 之一——否则拼成连续流后文档之间没有边界(python/dolma/tokenizer/executor.py:70-74 的硬校验)。


6. 关键细节与坑

  • 这是「轻混洗」不是「真混洗」:窗口只有 1 万篇、ring 只有 8 个文件;对混洗极敏感的训练(比如某些课程学习设置)需要在训练侧再补一层 shuffle。
  • EOS/BOS 二选一是硬要求:都不给会直接 ValueError(python/dolma/tokenizer/executor.py:70-74)——连续 token 流没有文档边界就没法用。
  • pad_token_id 缺省回退为 eos(python/dolma/tokenizer/executor.py:80-83):只警告不报错,语义上未必是你要的。
  • uint16 装不下大词表:词表 > 65535 时必须靠 dtype 自动选择或显式指定;写死 uint16 会静默溢出(inferred,Tokenizernp.min_scalar_type 自动选择能避免,但 CLI 显式传 uint16 时没有校验)。
  • 桶是按下标传的:grouped_source_prefixes 全量路径经 pickle 进每个子进程(python/dolma/tokenizer/executor.py:56-61 在子进程侧按下标取),文件数巨大时这份清单本身也是开销。
  • 看不出跨进程的全局去重/对齐:tokenize 阶段假设输入已经是干净的最终语料,不再做任何校验——垃圾进垃圾出。