@vanstriendaniel: datatrove — FineWeb、FineWeb2 和 FinePDFs 背后的数据处理库 — 刚刚发布了 0.10.0! - JobsPipelineExe…
摘要
Datatrove,FineWeb 数据集背后的数据处理库,发布了 0.10.0 版本,新增了 JobsPipelineExecutor 等功能,用于在 Hugging Face Jobs 上运行管道,并支持 HF 存储桶。
查看缓存全文
缓存时间: 2026/08/17 10:44
datatrove — 为 FineWeb、FineWeb2 和 FinePDFs 提供支持的数据处理库 — 正式发布 0.10.0 版本!
- JobsPipelineExecutor:在 Huggingface Jobs 上运行管道:支持扇出、多阶段依赖、重试和恢复功能,无需 Slurm 集群。
- HF 存储桶作为 DataFolder:直接读取、写入和记录至 hf://buckets/…
- 推理结果中保留推理输出
安装命令:
pip install datatrove[io,processing]
huggingface/datatrove
源码地址:https://github.com/huggingface/datatrove
DataTrove
DataTrove 是一个用于大规模处理、过滤和去重文本数据的库。它提供了一组预构建的常用处理模块,以及一个便于添加自定义功能的框架。
DataTrove 处理管道具有平台无关性,可直接在本地或 Slurm 集群上运行。其(相对)低的内存占用和多步骤设计使其非常适合处理大规模工作负载,例如处理 LLM 的训练数据。
通过 fsspec (https://filesystem-spec.readthedocs.io/en/latest/) 支持本地、远程及其他文件系统。
目录
安装
需要 Python 3.10+。
uv sync
可用版本(通过重复使用 --extra 进行组合,例如 uv sync --extra processing --extra s3):
all:安装所有依赖:uv sync --extra allio:用于读取warc/arc/wet文件和 arrow/parquet/Optimized-parquet (https://huggingface.co/docs/hub/en/datasets-libraries#optimized-parquet-files) 格式的依赖:uv sync --extra ioprocessing:用于文本提取、过滤和分词的依赖:uv sync --extra processings3:S3 支持:uv sync --extra s3cli:命令行工具:uv sync --extra cliray:分布式计算引擎:uv sync --extra rayinference:LLM 推理管道:uv sync --extra inferencedecont:使用 lighteval 进行去污染:uv sync --extra decontmultilingual:多语言文本处理:uv sync --extra multilingual
快速入门示例
你可以查看以下 示例:
- fineweb.py 完整复现 FineWeb 数据集 (https://huggingface.co/datasets/HuggingFaceFW/fineweb)
- process_common_crawl_dump.py 完整管道:读取 commoncrawl warc 文件,提取文本内容,进行过滤,并将结果数据保存到 s3。在 Slurm 上运行。
- tokenize_c4.py 直接从 Hugging Face Hub 读取数据,使用
gpt2分词器对 C4 数据集的英文部分进行分词。 - estimate_tokens.py 估算大型 HF 数据集的总 token 数——这在创建随机混洗子集(例如从数万亿 token 的数据集中抽取 100B token)时设置正确的
SamplerFilter比率时是必需的。它会流式处理每个数据集的一小部分样本,收敛到平均每个文档的 token 数,然后乘以总行数。 - smol_data.py 为多个大型 Hugging Face 数据集构建约 100B token 的子集、50-30-20 的混合数据以及混洗变体。
- minhash_deduplication.py 运行文本数据 MinHash 去重的完整管道。
- sentence_deduplication.py 运行句子级别精确去重的示例。
- exact_substrings.py 运行 ExactSubstr 的示例(需要此仓库 (https://github.com/google-research/deduplicate-text-datasets))。
- finephrase.py 使用多种提示模板大规模生成合成数据集的独立示例。
术语
pipeline:要执行的一系列处理步骤(读取数据、过滤、写入磁盘等)。executor:在给定的执行环境(slurm、多 CPU 机器等)上运行特定管道。job:在给定执行器上运行一个管道。task:一个job由多个task组成,用于并行化执行,通常每个task处理一个数据shard。Datatrove 会跟踪已完成的任务,当你重新启动时,只有未完成的任务会运行。file:单个输入文件(.json、.csv 等)。
请注意,每个文件将由单个
task处理。Datatrove 不会自动将文件拆分为多个部分,因此要实现完全并行化,你应该拥有多个中等大小的文件,而不是单个大文件。
shard:一组输入数据(通常是一组file),将被分配给特定的task。每个task将处理来自完整输入文件列表的不同且不重叠的数据shard。worker:一次执行单个任务的计算资源,例如,如果你有 50 个 CPU 核心,可以使用workers=50运行 LocalPipelineExecutor,同时执行 50 个task(每个 CPU 一个)。一旦一个worker完成一个task,它将开始处理另一个等待的task。
你的
task数量控制了可以并行化的程度,也决定了每个独立处理单元所需的时间。如果你的task数量较少(因此每个任务必须处理大量文件)并且它们失败了,你将不得不从头开始重新运行;而如果你的task数量较多(每个任务处理的文件较少),那么每个失败任务重新运行所需的时间就会少得多。 [!CAUTION] 如果你的task数量 >file数量,某些任务将不会处理任何数据,因此通常没有必要将task数量设置得比file数量更大。
示例
运行一个处理 10000 个 file 的 job,在拥有 100 个 CPU 核心(workers)的机器上。如果我们选择使用 1000 个 task,每个任务将处理一个包含 10 个文件的 shard。workers=100 意味着我们可以同时处理 100 个 task。
管道
DataTrove 文档
每个管道块都处理 datatrove Document 格式的数据:
text:每个样本的实际文本内容。id:此样本的唯一 ID(字符串)。metadata:一个字典,可存储任何附加信息。
管道块类型
每个管道块接受一个 Document 生成器作为输入,并返回另一个 Document 生成器。
- readers:从不同格式读取数据并生成
Document。 - writers:以不同格式将
Document保存到磁盘/云存储。 - extractors:从原始格式(如网页 HTML)中提取文本内容。
- filters:根据特定规则/标准过滤(移除)一些
Document。 - stats:收集数据集统计信息的块。
- tokens:对数据进行分词或计算 token 数的块。
- dedup:用于去重的块。
完整管道
管道被定义为一系列管道块。例如,以下管道将从磁盘读取数据,随机过滤(移除)一些文档,并将其写回磁盘:
from datatrove.pipeline.readers import CSVReader
from datatrove.pipeline.filters import SamplerFilter
from datatrove.pipeline.writers import JsonlWriter
pipeline = [
CSVReader(
data_folder="/my/input/path"
),
SamplerFilter(rate=0.5),
JsonlWriter(
output_folder="/my/output/path"
)
]
执行器
管道具有平台无关性,这意味着同一个管道可以在不同的执行环境中无缝运行,无需修改其步骤。每个环境都有自己的 PipelineExecutor。所有执行器通用的一些选项:
pipeline:由应运行的管道步骤组成的列表。logging_dir:一个 datafolder,用于保存日志文件、统计信息等。不要为不同的管道/作业重用文件夹,因为这会覆盖你的统计信息、日志和完成状态。skip_completed(布尔值,默认为True):Datatrove 会跟踪已完成的任务,以便在你重新启动作业时可以跳过它们。设置为False以禁用此行为。randomize_start_duration(整数,默认为0):延迟每个任务启动的最大秒数,以防止所有任务同时启动并可能过载系统。 调用执行器的run方法来执行其管道。
Datatrove 通过在
${logging_dir}/completions文件夹中创建标记(空文件)来跟踪哪些任务已成功完成。作业完成后,如果其某些任务失败,你可以简单地重新启动完全相同的执行器,Datatrove 将进行检查并仅运行之前未完成的任务。 [!CAUTION] 如果因为某些任务失败而重新启动管道,请不要更改任务总数,因为这会影响输入文件的分配/分片。
LocalPipelineExecutor
此执行器将在本地机器上启动管道。 选项:
tasks:要运行的任务总数。workers:同时运行的任务数。如果为-1,则无限制。任何> 1的值将使用多进程来执行任务。start_method:用于生成多进程 Pool 的方法。如果workers为 1,则忽略此选项。
示例执行器
from datatrove.executor import LocalPipelineExecutor
executor = LocalPipelineExecutor(
pipeline=[...],
logging_dir="logs/",
tasks=10,
workers=5
)
executor.run()
多节点并行
你可以通过使用 local_tasks 和 local_rank_offset 让不同的节点/机器处理总任务的不同部分。对于每个节点/实例/机器,使用以下选项启动:
tasks:要执行的总任务数(跨所有机器)。此值在每台机器上必须相同,否则输入文件的分配可能会重叠! 示例:500local_tasks:此特定机器上将执行的任务数。请注意,你可以为每台机器使用不同的值。示例:100local_rank_offset:此机器上要执行的第一个任务的排名。如果这是你启动作业的第 3 台机器,并且前 2 台机器分别运行了 250 和 150 个作业,那么当前机器的此值应为400。 要获得最终的合并统计信息,你必须在一个包含所有机器统计信息的路径上手动调用merge_stats脚本。
SlurmPipelineExecutor
此执行器将在 Slurm 集群上启动管道,使用 Slurm 作业数组来分组和管理任务。 选项:
tasks:要运行的任务总数。必需time:Slurm 时间限制字符串。必需partition:Slurm 分区。必需workers:同时运行的任务数。如果为-1,则无限制。Slurm 将一次运行workers个任务。(默认:-1)job_name:Slurm 作业名称(默认:“data_processing”)depends:另一个 SlurmPipelineExecutor 实例,将作为此管道的依赖项(当前管道将仅在依赖的管道成功完成后才开始执行)sbatch_args:包含你希望传递给 sbatch 的任何其他参数的字典。slurm_logs_folder:保存 slurm 日志文件的位置。如果logging_dir使用本地路径,它们将保存在logging_dir/slurm_logs下。否则,它们将作为当前目录的子目录保存。
其他选项
cpus_per_task:为每个任务分配的 CPU 数(默认:1)qos:Slurm QOS(默认:“normal”)mem_per_cpu_gb:每个 CPU 的内存,单位 GB(默认:2)env_command:自定义命令,用于在需要时激活 Python 环境。condaenv:要激活的 conda 环境。venv_path:要激活的 Python 环境路径。max_array_size:$ scontrol show config中的 MaxArraySize 值。如果任务数量超过此数字,它将分成多个数组作业(默认:1001)。max_array_launch_parallel:如果由于 max_array_size 需要多个作业,是否并行启动它们(默认:False)stagger_max_array_jobs:当 max_array_launch_parallel 为 True 时,这决定了在启动每个并行作业之间等待的秒数(默认:0)run_on_dependency_fail:当依赖的作业完成时开始执行,即使它已失败(默认:False)randomize_start:在约 3 分钟的窗口内随机化作业中每个任务的启动时间。当频繁访问 S3 存储桶时非常有用。(默认:False)
示例执行器
from datatrove.executor import SlurmPipelineExecutor
executor1 = SlurmPipelineExecutor(
pipeline=[...],
job_name="my_cool_job1",
logging_dir="logs/job1",
tasks=500,
workers=100, # 省略以一次运行所有
time="10:00:00", # 10 小时
partition="hopper-cpu"
)
executor2 = SlurmPipelineExecutor(
pipeline=[...],
job_name="my_cool_job2",
logging_dir="logs/job2",
tasks=1,
time="5:00:00", # 5 小时
partition="hopper-cpu",
depends=executor1 # 此管道将在 executor1 成功完成后才启动
)
# executor1.run()
executor2.run() # 这实际上会启动 executor1,因为它是依赖项,所以不需要显式启动它
RayPipelineExecutor
此执行器将在 Ray 集群上启动管道,使用 Ray 任务进行并行执行。 选项:
tasks:要运行的任务总数。workers:同时运行的任务数。如果为-1,则无限制。Ray 将一次运行workers个任务。(默认:-1)depends:另一个 RayPipelineExecutor 实例,将作为此管道的依赖项(当前管道将仅在依赖的管道成功完成后才开始执行)
其他选项
cpus_per_task:为每个任务分配的 CPU 数(默认:1)mem_per_cpu_gb:每个 CPU 的内存,单位 GB(默认:2)ray_remote_kwargs:传递给 ray.remote 装饰器的额外 kwargs
示例执行器
import ray
from datatrove.executor import RayPipelineExecutor
ray.init()
executor = RayPipelineExecutor(
pipeline=[...],
logging_dir="logs/",
tasks=500,
workers=100 # 省略以一次运行所有
)
executor.run()
JobsPipelineExecutor
实验性 — 可能随时更改或移除。
在 Hugging Face Jobs (https://huggingface.co/docs/huggingface_hub/en/guides/jobs) 上运行管道:与 Slurm 相同的协调模型,但每个任务块在云 Job 中运行。有关用法,请参阅 JobsPipelineExecutor 文档字符串 和 *_jobs.py 示例。
日志
对于具有 logging_dir mylogspath/exp1 的管道,将创建以下文件夹结构:
查看文件夹结构
└── mylogspath/exp1
├── executor.json ⟵ 执行器选项和管道步骤的 JSON 转储
├── launch_script.slurm ⟵ 用于启动此作业的 slurm 配置文件(如果运行在 Slurm 上)
相似文章
Datasette 1.0a40
Datasette 发布了 1.0a40 版本,包含安全修复、插件的新后台任务功能、迁移到 httpx2,以及各种错误修复,为稳定版 1.0 做准备。
datasette 1.0a39
Datasette 1.0a39 版本已发布,包含安全更新和数据探索增强功能。
datasette 1.0a37
Datasette 1.0a37 是 Simon Willison 宣布的一款用于探索和发布数据的开源工具的 alpha 版本。
datasette 1.0a28
Datasette 1.0a28 alpha 版本修复了前一个 alpha 版本中发现的兼容性错误和资源管理问题,包括修复 execute_write_fn() 回调、数据库清理方法,以及新增用于测试中自动清理的 pytest 插件。
datasette 1.0a29
Datasette 1.0a29 已发布,包含新的实用方法、对空表格的 UI 改进,以及在 Codex CLI 协助下修复的竞争条件等错误修复。