使用 Hugging Face 和 Dask 扩展 AI 数据处理规模
Hugging Face 和 Dask 提供了一个可扩展的框架,用于处理超出本地内存限制的大规模 AI 数据集。通过将 Dask 的分布式计算能力与 Hugging Face 的 transformers 和 datasets 相结合,用户可以将 AI 任务——例如模型推理和数据过滤——从笔记本电脑上的少量行扩展到跨多 GPU 云集群的数亿行。
使用 Dask 进行分布式数据处理
Dask 支持 out-of-core 计算,使得通过将数据集拆分为可管理的块来处理超出系统内存容量的数据成为可能。对于熟悉 pandas 的用户,Dask DataFrame 提供了类似的 API,简化了从本地原型到大规模生产的过渡。
使用 Dask 进行 AI 数据处理的主要优势包括:
- 高效加载: Dask 原生支持 Parquet,这是 Hugging Face 数据集的默认格式,能够实现高效的列式过滤和压缩。
- 并行执行:
map_partitions函数允许用户在更大的 Dask DataFrame 中的每个 pandas DataFrame 分区上并行应用自定义函数(例如模型推理)。 - 分布式写入: Dask 支持并行将结果写回 Parquet 格式,可与 Hugging Face 数据集仓库集成。
扩展模型推理:从 Pandas 到 Dask
为了演示扩展,Hugging Face 使用了 FineWeb 数据集——包含 15 万亿个英文网页数据的 token——以及 FineWeb-Edu 分类器来识别高教育价值的网页。
使用 Pandas 进行本地原型
在小规模(例如 100 行)时,可以使用 pandas 运行 FineWeb-Edu 分类器。在配备 GPU 的 M1 Mac 上,此过程大约需要 10 秒。工作流涉及使用 Hugging Face 的 pipeline 进行文本分类,在函数内部会动态选择硬件设备(CUDA、MPS 或 CPU),以确保代码后续分发时的兼容性。
扩展至 2.11 亿行
当扩展到 2.11 亿行的数据集(占用磁盘 432 GB 的爬取数据的一部分)时,串行处理变得极其缓慢。通过切换到 Dask DataFrame 并从 Hugging Face 懒加载数据,任务实现了并行化。
在此扩展工作流中,compute_scores 函数通过 map_partitions 应用。为优化性能,Hugging Face pipeline 中的 batch_size 被提升(例如提升至 768),以更好地利用 GPU 硬件。
云端多 GPU 并行推理
为了实现最大吞吐量,Dask 可以部署在云基础设施上。在示例中,使用 Coiled 通过 AWS g5.xlarge 实例(配备 NVIDIA A10 Tensor Core GPU)来 provision 一个包含 100 个 worker 的集群。
基础设施自动化
Coiled 自动化了多个关键的部署步骤:
- 虚拟机供应: 自动启动支持 GPU 的云 VM。
- 环境设置: 处理 NVIDIA 驱动和 CUDA 运行时的安装。
- 包同步: 将本地 Python 包和文件同步到云 worker,以确保环境一致性。
性能与利用率
处理 2.11 亿行数据大约耗时 5 小时。监控显示硬件利用率高效,GPU 中位利用率为 100%,GPU 可用 24 GB 内存中位使用量为 21.5 GB。
潜在的 AI 用例
这种分布式处理模式可用于除文本分类之外的各种大规模 AI 任务,包括:
- 基因组数据过滤: 从海量基因组数据集中挑选感兴趣的特定基因。
- 结构化数据提取: 使用大语言模型将非结构化文本转换为结构化数据集。
- 网页数据清洗: 对 Common Crawl 的大规模抓取数据进行清洗和过滤。
- 多模态推理: 使用多模态模型分析大规模音频、图像或视频数据集。