跳到主要内容

数据截至 (上游 commit a649de79c14a)

01 · Pipeline 抽象与执行器

这一章讲什么: DataTrove 的骨架——PipelineStep 生成器协议、executor 怎么把同一条链按 rank 复制 N 份、completions 空文件怎么做断点续跑,以及 Local / Slurm 两种 executor 的分发机制。读完你能回答:tasks=8000 的一行配置底下到底发生了什么。


1. 它要解决的小问题

处理几十亿文档,必然要并行。但「并行」拆开来其实是四件独立的事:

  • 把输入数据切成不重叠的份
  • 把同一份处理逻辑跑很多个副本(本机进程、集群任务);
  • 副本挂了之后只重跑没做完的份
  • 各副本的处理量/丢弃数汇总成一份报告

DataTrove 的答案是:前两件事归 executor,第三件归一个空文件约定,第四件归每个 block 自带的 stats 对象。处理逻辑本身(读、过滤、写)对此完全无感——它们只看到一个 rank 和一个 world_size


2. 思路:pipeline 就是一串生成器函数

先建立直觉。一个处理块的最小形态是什么?答案是:吃进一个文档生成器,吐出一个文档生成器

# 示意,非源码
def uppercase_everything(data, rank=0, world_size=1):
for doc in data: # doc 有 text / id / metadata
doc.text = doc.text.upper() # 改
yield doc # 放行;不 yield 就是丢弃

这个签名 (data, rank, world_size) -> generator 就是全部协议(src/datatrove/executor/base.py:130-132 对每个 callable 步骤就是这么调用的)。三个好处:

  • 组合免费writer(filter(reader())) 就是函数套函数,组装时一个文档都不读;
  • 内存 O(1):文档流过即处理,不落地,整条链的内存占用与数据集大小无关;
  • 分片内嵌rank / world_size 一路传到链上每个环节,谁需要切分(通常是 reader)谁自己用。

PipelineStepsrc/datatrove/pipeline/base.py:9)只是给这个协议加上工程外壳:抽象方法 run:98-114,注意默认实现是 yield from data——上游有数据就透传)、stat_update 记统计(:38-54)、track_time 计时(:81-93),以及一个巧妙点——__new__:22-32)在实例化时沿 MRO 收集 _requires_dependencies 并检查安装,所以 URLFilter 声明了 _requires_dependencies = ["tldextract", "fasteners", ("ahocorasick", "pyahocorasick")]src/datatrove/pipeline/filters/url_filter.py:46),缺依赖在构造时就报「请 pip install xxx」,而不是跑到一半才炸。tuple 形式 (module名, pip名) 用来处理模块名与包名不一致的情况(src/datatrove/utils/_import_utils.py:10-33)。


3. 图示:executor 的视角

launch (本地 shell / slurm 提交节点)


Executor.run()
│ ① get_incomplete_ranks() 列 completions/ 目录取差集

┌─ rank 0 ─┐ ┌─ rank 1 ─┐ ┌─ rank N-1 ─┐
│ reader │ │ reader │ ... │ reader │ 每个 rank 是同一条
│ ↓ │ │ ↓ │ │ ↓ │ pipeline 的一份拷贝
│ filter │ │ filter │ │ filter │ (reader 按 rank 取
│ ↓ │ │ ↓ │ │ ↓ │ 不同文件 shard)
│ writer │ │ writer │ │ writer │
└────┬─────┘ └────┬─────┘ └─────┬──────┘
▼ ▼ ▼
stats/00000.json + completions/00000 (每 rank 各写各的)

怎么读这张图: executor 不关心链上有什么 block;它只看到 range(world_size) 里的 rank,负责为每个未完成的 rank 调用一次 _run_for_rank,然后记账。


4. 真实实现:_run_for_rank 的十六行核心

PipelineExecutor._run_for_ranksrc/datatrove/executor/base.py:102)是全部 executor 共用的心脏。核心三步:

# src/datatrove/executor/base.py:129-138
pipelined_data = None
for pipeline_step in self.pipeline:
if callable(pipeline_step):
pipelined_data = pipeline_step(pipelined_data, rank, self.world_size)
elif isinstance(pipeline_step, Sequence) and not isinstance(pipeline_step, str):
pipelined_data = pipeline_step # 直接给一个 list[Document] 也行
else:
raise ValueError
if pipelined_data:
deque(pipelined_data, maxlen=0) # 点火:消费到底

三点值得停下来看:

  1. 组装与执行分离。 循环只是把生成器一层层套起来,真正的数据流动由 deque(pipelined_data, maxlen=0) 触发——这是「消费一个迭代器并丢弃结果」的标准技法(maxlen=0 的双端队列不存任何元素)。
  2. 记账在成功之后。 链跑完后才写 stats/{rank:05d}.jsonmark_rank_as_completed(rank):143-148),而 mark_rank_as_completed 的实现就是创建一个空文件 completions/{rank:05d}:167-177)。任务中途崩了,completions 文件不存在,下次重跑自动补上。
  3. 入口的短路。 方法开头 is_rank_completed(rank) 命中就直接返回空 stats(:116-118);skip_completed=False 可以关掉这个行为(:156-165)。

这个「空文件 = 完成」的约定是整个续跑机制的全部。get_incomplete_ranks:179-195)就是列 completions/ 目录文件名、和 range(world_size) 做差集。朴素,但对「数据处理任务要么全对要么重跑」的语义完全够用——重跑同一个 rank 会覆盖自己写的那几个输出文件(文件名含 ${rank},见第 2 章),天然幂等。


5. LocalPipelineExecutor:本机多进程

LocalPipelineExecutorsrc/datatrove/executor/local.py:15)把上面那套语义套在本机上。三个参数决定形态:

  • tasks:总份数(= world_size:169-176);
  • workers:同时跑几份,workers=1 就退化成顺序循环(:125-130),大于 1 用 multiprocess 进程池 + imap_unordered:134-146);
  • local_tasks + local_rank_offset:多机分摊同一个 job 时,本机只跑 [offset, offset+local_tasks) 这段 rank(:36-59),任务分发仍是全局确定的。

两个细节:

  • 每份任务拿到 deepcopy 的 pipeline(顺序分支里 self.pipeline = deepcopy(pipeline):129),避免多 rank 间 stats 互相污染;
  • depends=:如果声明了依赖另一个本地 executor,先递归启动它,然后每 2 分钟轮询 get_incomplete_ranks,等依赖的 rank 全部完成才开始自己(:96-105)。多个 job 串成 DAG 就靠这个。

跑完后所有 rank 的 stats 用 sum(stats, start=PipelineStats()) 合并写 stats.json:147-151)——合并靠 PipelineStats.__add__src/datatrove/utils/stats.py:140-143),底层的 MetricStats.__add__:249-274)用并行方差公式合并均值/方差,而单个任务的均值方差本来就是用 Welford 在线算法算的(MetricStats.update:217-239)。所以「8000 个任务每步各扔掉多少」这份报告不需要任何集中式收集,纯靠文件落盘 + 离线合并。


6. SlurmPipelineExecutor:把整个 executor pickle 进 job array

Slurm 版是生产主力(FineWeb 就用它),机制是「自我序列化 + job array」。分两半看。

6.1 提交侧(登录节点上)

run() 先查环境变量:没有 SLURM_ARRAY_TASK_ID 说明自己在提交节点上,走 launch_job()src/datatrove/executor/slurm.py:199-227)。launch_job:263)依次做:

  1. 先发射依赖depends 的 executor 若还没 job_id 就先递归 launch_job(),然后把依赖换成 depends_job_id,并 self.depends = None——注释写明是「避免把整条依赖链 pickle 进去」(:272-279);
  2. 算未完成的 rank,全部完成则 job_id = -1 直接收工(:281-285);
  3. 把自己 dill 成 executor.pik:287-291),同时把 ranks_to_run 存成 JSON——「只存一次,避免竞态」(:294-296);
  4. 生成 sbatch 脚本get_launch_file_contents:366-402)+ sbatch 参数(get_sbatch_args:332-364,含 --array=0-{max_array-1}%{workers}、requeue、邮件等),调用 sbatch 提交(launch_slurm_job:420-435);任务数超集群 max_array_size 时拆成多个 array job,靠 RUN_OFFSET 环境变量错开 rank 区间(:320-328);
  5. 顺手提交一个 merge_stats 收尾 job,afterok 依赖主 job,把各任务 stats 合成总表(launch_merge_stats:229-247)。

6.2 计算侧(集群节点上)

job array 的每个元素启动后跑的是 srun ... launch_pickled_pipeline executor.pik:307-309)。这个命令对应 src/datatrove/tools/launch_pickled_pipeline.py:14-18,全部逻辑就三行:dill 加载 executor、调 executor.run()

这次 run() 里有 SLURM_ARRAY_TASK_ID 了,于是走另一半:把 array index 映射回本 executor 的 rank 区间(乘以 tasks_per_job:201-204),从 ranks_to_run.json 读出真正的 rank 列表(:205-206),逐个调 _run_for_rank:219-224)。

提交节点 集群
───────── ─────
launch_job()
dill → executor.pik ──► sbatch --array=0-N
ranks_to_run.json │ 每个 array 元素:

launch_pickled_pipeline
dill.load → run()
│ SLURM_ARRAY_TASK_ID → rank 区间

_run_for_rank(rank) × tasks_per_job

还有两个集群生存细节:注册 SIGUSR1 等信号的 handler,收到就 scontrol requeue 把自己重新排队再退出(requeue_handler:41-45)——抢占式集群上实例被回收时能自动重投;randomize_start_duration 让每个任务随机睡 0~N 秒再开工(executor/base.py:125-126),避免几千个任务同时 list S3 bucket 把对象存储打爆(examples/fineweb.py:69 给了 180 秒)。

6.3 Ray 与 HF Jobs

RayPipelineExecutorsrc/datatrove/executor/ray.py:385)用 Ray actor(RankWorker:29)替掉进程池,思路相同:actor 里调 _run_for_rank,支持跨节点。JobsPipelineExecutor 把同一模型搬到 HF Jobs 云上。两者不改变任何 pipeline 写法——这正是「pipeline 与平台无关」(README.md:134)的兑现方式。


7. 关键细节 / 坑

  • tasks 数定了就别改。 shard 划分是 all_files[rank::world_size]src/datatrove/io.py:164-180),改了 tasks 就改了文件→rank 的映射,续跑会错乱。README 用 CAUTION 标了这条(README.md:148-149)。
  • tasks > 文件数没有意义:多出来的 rank 分到空 shard,白排队(README.md:93-94)。反过来,少量大文件会卡死并行度——切分单位是文件不是字节。
  • logging_dir 不能复用:stats/logs/completions 都写在里面,两个不同 pipeline 共用一个目录会互相覆盖、互相误判「已完成」(README.md:138)。
  • 本地多进程的 start_method 默认 forkserverexecutor/local.py:44)。block 里在模块级持有不可 pickle 的资源会在 fork 后出问题;README 建议自定义函数把 import 写进函数体(README.md:621-622)。
  • rank 与「文档在文件里的顺序」是隐性契约。 去重 stage 产出的 .remove 清单按 (rank, doc_idx) 定位文档(见第 4 章),所以 reader 前面开 shuffle_files=True 会让去重对错目标——BaseDiskReader 的 docstring 明令「do not use with dedup blocks」(src/datatrove/pipeline/readers/base.py:132)。
  • 代码里看不出的:Slurm 侧对「部分任务失败后整 job 的状态」依赖集群本身的 array 语义,run_on_dependency_fail 决定依赖失败时是否照跑(slurm.py:98:258);具体重试几次、退避如何,是集群配置而非库代码控制的。

8. 本章代码地图

主题文件路径符号名
处理块协议src/datatrove/pipeline/base.pyPipelineStep.runPipelineStep.__new__stat_update
executor 心脏src/datatrove/executor/base.pyPipelineExecutor._run_for_rankmark_rank_as_completedget_incomplete_ranks
本地执行src/datatrove/executor/local.pyLocalPipelineExecutor.run_launch_run_for_rank
Slurm 提交src/datatrove/executor/slurm.pySlurmPipelineExecutor.launch_jobget_launch_file_contentsrequeue_handler
pickle 入口src/datatrove/tools/launch_pickled_pipeline.pymain
统计与合并src/datatrove/utils/stats.pyMetricStats.updateMetricStats.__add__PipelineStats
依赖惰性检查src/datatrove/utils/_import_utils.pycheck_required_dependencies