@vanstriendaniel: datatrove — FineWeb、FineWeb2 和 FinePDFs 背后的数据处理库 — 刚刚发布了 0.10.0! - JobsPipelineExe…

X AI KOLs Following 工具

摘要

Datatrove,FineWeb 数据集背后的数据处理库,发布了 0.10.0 版本,新增了 JobsPipelineExecutor 等功能,用于在 Hugging Face Jobs 上运行管道,并支持 HF 存储桶。

datatrove — FineWeb、FineWeb2 和 FinePDFs 背后的数据处理库 — 刚刚发布了 0.10.0 版本! - JobsPipelineExecutor:在 @huggingface Jobs 上运行管道:支持扇出、多阶段依赖、重试和恢复。无需 Slurm 集群。 - HF 存储桶作为 DataFolder 使用:直接读取、写入和记录到 hf://buckets/... - 推理结果中保留了推理输出。 pip install datatrove[io,processing]
查看原文
查看缓存全文

缓存时间: 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 all
  • io:用于读取 warc/arc/wet 文件和 arrow/parquet/Optimized-parquet (https://huggingface.co/docs/hub/en/datasets-libraries#optimized-parquet-files) 格式的依赖:uv sync --extra io
  • processing:用于文本提取、过滤和分词的依赖:uv sync --extra processing
  • s3:S3 支持:uv sync --extra s3
  • cli:命令行工具:uv sync --extra cli
  • ray:分布式计算引擎:uv sync --extra ray
  • inference:LLM 推理管道:uv sync --extra inference
  • decont:使用 lighteval 进行去污染:uv sync --extra decont
  • multilingual:多语言文本处理: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 数量更大。

示例

运行一个处理 10000filejob,在拥有 100 个 CPU 核心(workers)的机器上。如果我们选择使用 1000task,每个任务将处理一个包含 10 个文件的 shardworkers=100 意味着我们可以同时处理 100task

管道

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_taskslocal_rank_offset 让不同的节点/机器处理总任务的不同部分。对于每个节点/实例/机器,使用以下选项启动:

  • tasks:要执行的总任务数(跨所有机器)。此值在每台机器上必须相同,否则输入文件的分配可能会重叠! 示例:500
  • local_tasks:此特定机器上将执行的任务数。请注意,你可以为每台机器使用不同的值。示例:100
  • local_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

Simon Willison's Blog

Datasette 发布了 1.0a40 版本,包含安全修复、插件的新后台任务功能、迁移到 httpx2,以及各种错误修复,为稳定版 1.0 做准备。

datasette 1.0a39

Simon Willison's Blog

Datasette 1.0a39 版本已发布,包含安全更新和数据探索增强功能。

datasette 1.0a37

Simon Willison's Blog

Datasette 1.0a37 是 Simon Willison 宣布的一款用于探索和发布数据的开源工具的 alpha 版本。

datasette 1.0a28

Simon Willison's Blog

Datasette 1.0a28 alpha 版本修复了前一个 alpha 版本中发现的兼容性错误和资源管理问题,包括修复 execute_write_fn() 回调、数据库清理方法,以及新增用于测试中自动清理的 pytest 插件。

datasette 1.0a29

Simon Willison's Blog

Datasette 1.0a29 已发布,包含新的实用方法、对空表格的 UI 改进,以及在 Codex CLI 协助下修复的竞争条件等错误修复。