用 Dask 在 Hugging Face Jobs 上生成合成训练数据

0 阅读

合成数据为什么需要专门的流水线

训练大模型时,数据往往比算法更关键。但现实中,高质量标注数据要么难获取,要么成本太高。于是很多人转向合成数据——让语言模型自己生成训练样本。

听起来简单:给模型一段原文,让它输出问题和答案。但真要跑起来,你会发现这其实是个混合型计算任务:

  • 准备阶段(读原始数据、抽样、分批)全是 CPU 活;
  • 生成阶段必须用 GPU 跑模型;
  • 后处理阶段(解析、验证、去重、写文件)又回到 CPU。

如果全塞进一个 GPU 进程里跑,等于让昂贵的显卡干搬砖的活。分开写多个脚本?又得引入额外的调度系统,中间结果还得存到磁盘,拖慢整体速度。

理想方案是:用一个统一的任务图,把 CPU 和 GPU 工作串起来,按需分配资源。Dask 正好擅长这事,而 Hugging Face Jobs 提供了靠近数据的异构算力。缺的只是一个能把两者连起来的胶水——这就是 hfdask 出现的原因。

hfdask 是什么?

hfdask 是个轻量级工具,作用就一个:根据 YAML 配置,在 Hugging Face Jobs 上临时拉起一个 Dask 集群

它会启动两类 Job:

  • 协调器(coordinator):跑 Dask 调度器 + 用户主脚本;
  • 工作节点(workers):可以是纯 CPU,也可以是带 GPU 的机器。

所有节点通过加密的 Iroh 网络互联,不暴露公网端口。每个 GPU 节点会自动给第一个 worker 分配一块显卡,并打上 GPU 资源标签。你的代码只需用标准 Dask API,比如 resources={"GPU": 1},就能把任务精准调度到 GPU 上。

最关键的是:你的业务逻辑完全不用改。hfdask 只管集群生命周期——提交 Job、装依赖、挂载数据、清理资源——剩下的还是你熟悉的 Dask 编程。

构建问答合成流水线

我们以生成新闻领域的问答对为例,完整走一遍流程。目标是从 AG News 数据集的新闻正文里,让模型生成“问题-答案-支持引文”三元组。

数据源:AG News 测试集

选用 fancyzhx/ag_news 的 test split。每条记录包含新闻文本和类别标签(World/Sports/Business/Sci/Tech)。我们保留原始行号作为溯源 ID。

cluster.yaml 里直接挂载数据集:

mounts:
  - source: hf://datasets/fancyzhx/ag_news
    revision: eb185aade064a813bc0b7f42de02595523103ca4
    target: /dataset

CPU worker 读取 Parquet 文件,做三件事:

  1. 按类别均衡采样(每类 32 条);
  2. 把整数标签转成可读字符串;
  3. 每 16 条切一个分区,变成 Dask 任务单元。

这样既保证类别平衡,又避免首次运行太耗资源。

生成模型:Qwen3-0.6B

为了演示方便,选了小模型 Qwen/Qwen3-0.6B,能在单张 L4 显卡上跑起来。生产环境当然要用更大的模型。

同样通过挂载引入:

mounts:
  - source: hf://models/Qwen/Qwen3-0.6B
    revision: c1899de289a04d12100db370d81485cdf75e47ca
    target: /model

重点来了:不能每次生成都重新加载模型。我们用 Dask 的 WorkerPlugin 机制,在 GPU worker 启动时初始化一次 vLLM 引擎,并挂在 worker 对象上供后续任务复用。

class GeneratorSetup(WorkerPlugin):
    def setup(self, worker: Worker) -> None:
        if not worker.state.total_resources.get("GPU", 0):
            return  # CPU worker 跳过
        worker.generator = LLM(model="/model", ...)
        worker.sampling_params = SamplingParams(...)

这样,每个 GPU worker 有且仅有一个模型实例,避免重复加载开销。

输出格式:严格约束的 JSON

提示词明确要求模型返回特定结构的 JSON:

{
  "question": "...",
  "answer": "...",
  "supporting_quote": "原文中的 exact quote"
}

并用 Pydantic 做强校验:

  • 不允许多余字段;
  • 所有字段非空;
  • 支持引文必须原样出现在原文中。

这种“宁可拒收也不修复”的策略,能防止脏数据混进训练集。

集群配置:CPU + GPU 混合

项目用 uv 管理依赖,hfdask 放在 deploy 分组里:

uv add dask distributed pandas pyarrow pydantic "vllm==0.29.0; sys_platform == 'linux'"
uv add --group deploy "hfdask==0.1.1"

cluster.yaml 定义了一个 CPU 协调器 + 一个 L4 GPU worker:

coordinator:
  flavor: cpu-basic
  worker: true  # 协调器也兼做 CPU worker

workers:
  flavor: l4x1
  count: 1

mounts:
  - source: hf://models/Qwen/Qwen3-0.6B
    target: /model
  - source: hf://datasets/fancyzhx/ag_news
    target: /dataset
  - source: hf://buckets/your-namespace/synthetic-data/run-001
    target: /output
    read_only: false

注意三点:

  1. 协调器开启 worker: true,能分担 CPU 任务;
  2. 输出桶挂载为可写,用于存结果;
  3. 使用 vLLM 官方镜像,确保 CUDA 环境一致。

验证:不只是格式检查

通过 Pydantic 解析只是第一步。我们还加了业务规则:

  • 问题不超过 240 字符;
  • 答案不超过 1000 字符;
  • 支持引文必须 exact match 原文。

不符合的记录会被归入 rejected.parquet,并附上拒绝原因(如 supporting_quote_not_found)。这些“失败案例”对迭代提示词特别有用。

最后,对通过验证的问题做归一化去重(转小写、压缩空格),避免模型反复生成相似问题。

组装完整任务图

主函数用 Dask 注解明确指定任务位置:

# 1. CPU 准备数据
with dask.annotate(workers=cpu_workers):
    source_parts = [dask.delayed(prepare_source)(DATASET)]

![hfdask-architecture](https://cdn-uploads.huggingface.co/production/uploads/63fccb965d2bea4588be8fb1/I_LS35NX-cjzv08PvoS3l.png)

# 2. GPU 生成
with dask.annotate(resources={"GPU": 1}):
    generated_parts = [dask.delayed(generate_partition)(p) for p in source_parts]

# 3. CPU 验证+写入
with dask.annotate(workers=cpu_workers):
    validated = [dask.delayed(validate_partition, nout=2)(p) for p in generated_parts]
    summary = dask.delayed(write_dataset)(accepted_parts, rejected_parts, OUTPUT)

关键参数 optimize_graph=False 防止 Dask 自作聪明地合并不同资源类型的任务。

如何运行

先登录 Hugging Face,再执行:

uv run --group deploy hfdask run \
  --cluster cluster.yaml \
  generate_dataset.py

hfdask 会自动:

  1. 打包当前 Git 目录(不含 .gitignore 内容);
  2. 上传到私有 jobs-artifacts 桶;
  3. 提交协调器和 worker Job;
  4. 各 Job 启动后,用 uv 安装锁定依赖;
  5. 建立加密隧道,形成集群;
  6. 协调器跑主脚本,完成后自动清理所有 Job。

本地机器甚至不需要装 vLLM——依赖只在 Linux Job 里安装。

扩展建议

这套架构天然支持横向扩展:

  • 增加 GPU worker 数量workers.count),每个自带独立模型副本;
  • 调大输入分区数,确保 GPU 不闲着;
  • 监控“有效产出/美元”,而非单纯看吞吐量;
  • 长任务要分片写入,支持失败重试时跳过已完成部分。

但别轻易搞模型并行——那需要换镜像、改启动参数,复杂度陡增。多数场景下,多副本更简单可靠。

hfdask 的边界在哪里

hfdask 只负责“怎么跑”,不干涉“跑什么”:

hfdask 管 用户代码管
Job 提交与清理 提示词设计
环境依赖安装 生成质量
数据集挂载 输出 schema
加密通信 验证规则
资源发现 去重逻辑

这种分离让同一套 Dask 代码既能跑在本地,也能无缝迁移到 Hugging Face Jobs。

小结

合成数据生成不是简单的 API 调用,而是一个涉及多种计算资源的流水线。hfdask 填补了 Dask 与 Hugging Face Jobs 之间的空白,让你用熟悉的 Python 代码,高效利用异构硬件。整个过程无需改动业务逻辑,只需一个 YAML 配置,就能把 CPU 准备、GPU 生成、CPU 后处理串成一条流水线。

项目已开源,安装命令:

uv add --group deploy hfdask

无论是生成问答对、指令数据,还是做模型蒸馏,这个模式都值得一试。