Data-Juicer 分布式数据处理:基于 Ray 的大规模数据流水线实战指南
Data-Juicer 分布式数据处理:基于 Ray 的大规模数据流水线实战指南
导读
Data-Juicer 在单机模式下提供了丰富的数据处理算子,而面对百亿级样本、TB 级数据集的清洗与去重,单机模式无论在内存还是算力上都难以为继。本文以 docs/Distributed_ZH.md 为骨架,系统讲解 Data-Juicer 如何借助 Ray 与阿里云 PAI 实现大规模分布式数据处理:从引擎无关的算子设计、数据子集分割、JSON 流式读取补丁、分布式 MinHash-LSH 去重,到 RayAnalyzer 分布式数据分析的完整原理,并结合仓库内真实配置与源码,给出可直接运行的快速开始指南。读完本文,你将掌握如何用一条 dj-process/dj-analyze 命令把单机数据处理流水线无缝迁移到 Ray 集群,并理解大规模去重与数据切分的底层优化思路。
概览:为什么需要分布式数据处理
Data-Juicer 的分布式能力建立在 Ray(开源分布式计算框架)与阿里巴巴 PAI(机器学习平台)之上。其核心设计目标非常明确:几乎所有在单机模式下实现的算子都可以无缝运行在 Ray 分布式模式下。
这一目标之所以能达成,关键在于引擎无关的算子架构。为了让大规模场景下表现更优,项目还做了多项针对计算引擎的优化:
- 数据子集分割策略:用于平衡文件数目与进程数目,避免极端场景下的网络通信开销;
- 针对 Ray 与 Apache Arrow 的 JSON 文件流式 I/O 补丁:缓解 TB 级 JSONL 文件加载时的内存压力。
从官方文档给出的性能参考来看(以下数据均来自 docs/Distributed_ZH.md,为项目作者在 25~100 个阿里云节点上实测的结果):
- 在 6,400 个 CPU 核上处理包含 700 亿条样本的数据集,耗时仅约 2 小时;
- 在 3,200 个 CPU 核上处理包含 70 亿条样本的数据集,耗时约 0.45 小时;
- 在 1,280 个 CPU 核(8 节点集群)上,对 TB 级别数据集 执行 MinHash-LSH 去重,仅需约 3 小时。
这些数据反映的是特定集群配置与特定数据集上的实验结果,用于说明架构的可扩展性,不构成任何通用性能承诺。
更详细的设计与实验请参考论文《Data-Juicer 2.0: Cloud-Scale Adaptive Data Processing for Foundation Models》(arxiv.org/abs/2501.14755)。
实现与优化:Ray 模式下的四个关键技术点
1. 引擎无关的算子 + RayDataset / RayExecutor
Data-Juicer 的大部分算子的核心处理函数是引擎无关的——它们只面向数据本身编写,与底层执行引擎解耦。与 Ray 引擎的具体互操作被封装在两个类中:
- RayDataset:
DJDataset的子类,负责数据的分布式加载、map_batches算子分发、统计列维护与过滤; - RayExecutor:
BaseExecutor的子类,负责初始化 Ray、装配算子、执行流水线并导出结果。
两者都同时支持 Ray 的 Tasks(无状态任务)与 Actors(有状态长期驻留进程)两种执行原语。在 ray_dataset.py 的 _run_single_op 中可以看到:Mapper 与 Filter 算子默认通过 TaskPoolStrategy 以任务方式执行;当算子声明 use_ray_actor() 时,则切换为 ActorPoolStrategy(支持固定或弹性并发区间 <a href="https://link.gitcode.com/i/cc3681d5f3f5cb96521f52ac129b6c95" target="_blank">min_size, max_size]),并把算子类本身作为 Actor 构造传入 map_batches(见 _build_actor_pool_strategy,[ray_dataset.py)。
一个重要的约束体现在 ray_executor.py 的类文档中:
- 只支持
_supported_exec_modes中包含ray的算子;- 只支持加载 JSON 格式文件;
- 高级功能(如 checkpoint)不支持。
如果某个算子不支持 Ray 模式,执行时会抛出 NotImplementedError 并明确提示其支持的执行模式(见 ray_dataset.py)。
去重算子是一个例外:去重在单机模式下很难规模化(需要在全量数据上做全局等价类合并),因此项目提供了专门的 Ray 优化版本算子,并以特殊前缀 ray_ 开头,位于 data_juicer/ops/deduplicator/ 目录,例如 ray_bts_minhash_deduplicator、ray_document_deduplicator、ray_image_deduplicator、ray_video_deduplicator 等。
2. 数据子集分割:避免“文件太少、节点太多”的通信风暴
当在上万个节点上处理只有若干个文件的数据集时,Ray 会根据可用资源自动分割数据集文件并分发到所有节点,这可能带来极大的网络通信开销并拉低 CPU 利用率(相关细节可参考 Ray 的 autodetect_parallelism 机制与输出块调优文档)。
为了规避这种低效的默认执行计划,Data-Juicer 结合 Ray 与 Arrow 的特性,提前将原始数据集自动拆分为较小的文件。自动拆分策略有两个硬性约束:
- 单个文件大小设置为 128MB;
- 拆分后的子文件数量至少是集群可用 CPU 核心总数的两倍。
对应工具位于 tools/data_resplit.py。从源码看,其拆分逻辑非常务实:
- 通过
ray.cluster_resources().get("CPU", 1)获取集群 CPU 总数; - 计算
max_size = max(1, min(128, total_size / cpu_num / 4))(单位 MB),即目标文件大小会随集群规模自动收缩,保证子文件数量足够摊满所有核; - 使用
ray.data.from_pandas+map_batches并行执行split_jsonl,把大 JSONL 文件按字节阈值切分为{文件名}_{index}.jsonl子文件。
使用方式(命令行参数见脚本 main):
python tools/data_resplit.py --data-dir <数据集目录> --resplit-dir <切分输出目录> [--ray-address auto]
用户遇到此类性能问题时,既可以使用该工具,也可以根据自己的偏好自行拆分数据集。
3. JSON 文件的流式读取:绕开 Arrow 的原生缺失
基础模型数据集中大量数据以 JSONL 格式存储且体积巨大,流式读取 JSON 是常见刚需。但截至 Ray 2.40 / Arrow 18.1.0,Ray Datasets 底层并不原生支持流式读取 JSON(根因在 Apache Arrow 库)。为此 Data-Juicer :
- 开发了一个流式加载接口;
- 为 Apache Arrow 贡献了第三方补丁(对应 PR 见 data-juicer 的 PR #515 与 apache/arrow 的 PR #45084),以缓解内存不足问题。
在 ray_dataset.py 中可以看到该接口的落地实现:
RayDataset.read_json会转调自定义的read_json_stream(见 ray_dataset.py);read_json_stream优先使用 PyArrow 20.0.0+ 的pyarrow.json.open_json按 batch 流式读取,并处理跨 batch 的 schema 演化(pyarrow.unify_schemas+cast);对旧版本 PyArrow 则优雅降级到ray.data.read_json(见 ray_dataset.py);- 自定义
JSONStreamDatasource继承自 Ray 的JSONDatasource(或ArrowJSONDatasource,视 Ray 版本而定),并覆写_read_stream实现流式迭代。
使用该补丁后,Data-Juicer 的 Ray 模式会默认使用流式加载接口加载 JSON 文件;而当输入为 CSV 和 Parquet 文件时,Ray 模式下的流式读取也已经自动开启。
4. 分布式去重:MinHash-LSH + 多进程并查集 + BTS 负载均衡
Ray 模式下提供的优化版去重算子核心是 MinHash-LSH 技术,实现在 ray_bts_minhash_deduplicator.py(算子名 ray_bts_minhash_deduplicator)。其关键设计包括:
- 多进程并查集(Union-Find):用 Ray Actors 实现了分布式的
BTSUnionFind(见 ray_bts_minhash_deduplicator.py),每个 Actor 维护自己的哈希表与父节点映射; - 负载均衡的 BTS 算法:论文源自 IEEE 收录的 BTS Union-Find 工作。算法通过
edge_redistribution、balanced_union_find、communication、rebalancing、squeeze等阶段,在多个 Actor 之间按uid // BATCH_SIZE % parallel_num均匀分发边数据,迭代合并等价类直到收敛; - 远程 EdgeBuffer Actor:作为各 Union-Find Actor 之间交换边的中间缓冲区,用
ray.wait控制并发任务数量(max_pending_edge_buffer_task、num_edge_buffer_task_returns),防止任务堆积; - 惰性初始化:
_ensure_actors在集群完成自动扩容后才按“CPU 核数的一半”创建 Actor 池(union_find_parallel_num = max(1, int(ray.cluster_resources().get("CPU", 1) / 2)),见 ray_bts_minhash_deduplicator.py); - GPU 加速可选:
GPUMinHashActor借助 cuDF(cudf.Series.str.minhash)在 GPU 上计算 MinHash 签名,但仅支持 character 分词(见 ray_bts_minhash_deduplicator.py); - 增量去重变体:同文件还注册了
ray_bts_minhash_deduplicator_with_uid,要求输入含唯一__dj__uid列,可减少中间落盘 I/O 并支持对“已去重数据集 + 新数据”的增量联合去重(优先级按 uid 大小裁决)。
官方消融实验表明:相比该去重算子的初始实现版本,这些专门优化可带来 2~3 倍 提速(结果来自 docs/Distributed_ZH.md)。
性能结果:来自官方文档的实测参考
以下数据全部引自 docs/Distributed_ZH.md:
不同数据规模的数据处理
作者准备了一个 56 万条样本的多模态数据集,并用 1~125,000 倍的倍数扩展出不同规模的数据集,在十亿样本规模上验证了 Data-Juicer 的高扩展性(扩展后样本规模可达数百亿条)。
大规模数据集分布式去重
在 200GB、1TB、5TB 数据集上测试基于 MinHash 的 Ray 去重算子,CPU 核数从 640 核到 1280 核。结果规律为:数据集大小增长 5 倍,处理时间约增长 4.02~5.62 倍;CPU 核数翻倍,处理时间减少约 58.9%~67.1%。
| CPU 核数 | 200GB 耗时 | 1TB 耗时 | 5TB 耗时 |
|---|---|---|---|
| 4 × 160(640 核) | 11.13 分钟 | 50.83 分钟 | 285.43 分钟 |
| 8 × 160(1280 核) | 7.47 分钟 | 30.08 分钟 | 168.10 分钟 |
分布式数据分析:RayAnalyzer
除分布式数据处理外,Data-Juicer 还通过 RayAnalyzer 支持分布式数据分析,它是本地 Analyzer 的分布式版本。
工作原理
RayAnalyzer 复用本地 Analyzer 的 dj-analyze 命令行入口:当配置文件中 executor_type 设置为 ray 时,dj-analyze 会自动调度到 RayAnalyzer。工作流程如下:
- 加载数据 —— 通过 Ray 的分布式数据加载能力(支持 JSON、CSV、Parquet);
- 计算统计信息 —— 通过 Ray
map_batches分布式执行每个 Filter 的compute_stats,与本地模式使用相同的统计计算逻辑(对应 ray_dataset.py 中Filter分支的op.compute_stats调用,见 ray_dataset.py); - 聚合总体统计 —— 使用 Ray 原生聚合算子(
Count、Mean、Std、Min、Max)计算总体统计信息,无需 pandas 物化,适用于任意规模的数据集(见 ray_analyzer.py 中_compute_overall对ray.data.aggregate的使用,ray_analyzer.py); - 导出 —— 通过
RayExporter将带有统计列的数据集导出到磁盘。
其中 _flatten_stats_batch 会把嵌套的 stats/meta 字典列展开为顶层列,再按字段类型区分数值列与列表数值列分别聚合,最终结果写入工作目录 analysis/overall.csv 与 analysis/overall.md。
与本地 Analyzer 的差异
| 特性 | 本地 Analyzer | RayAnalyzer |
|---|---|---|
| 总体统计(count/mean/std/min/max) | 支持 | 支持 |
| 逐列分布图表 | 支持 | 不支持 |
| 相关性分析 | 支持 | 不支持 |
| 分位数 | 支持 | 不支持 |
| 可扩展性 | 单机 | 分布式(Ray 集群) |
因此,RayAnalyzer 聚焦于在大规模数据上计算总体统计信息;如果还需要详细的可视化与相关性分析,建议对采样子集使用本地 Analyzer。
快速开始:三个可运行的官方示例
环境准备:安装依赖并启动 Ray 集群
首先安装 Data-Juicer 及 dist 扩展依赖(distributed 组包含 Ray 及其他分布式相关依赖库):
uv pip install -v -e . # 安装 Data-Juicer 的最小依赖需求
uv pip install -v -e ".[distributed]" # 包括 Ray 以及其他分布式相关的依赖库
然后启动一个 Ray 集群(完整说明可参考 Ray 官方文档 starting-ray):
# 启动一个集群并作为头节点
ray start --head
# (可选)在其他节点或机器上连接集群
ray start --address='{head_ip}:6379'
[!Important] 如果要在多个节点上运行示例,需要把示例数据集放到**共享磁盘(如 NAS)**上,并将结果数据集导出到那里——只需修改配置文件中的
dataset_path与export_path参数即可。
项目在 demos/process_on_ray/ 目录下准备了简单示例,包括 2 个配置文件和 2 个测试数据集:
demos/process_on_ray
├── configs
│ ├── demo.yaml
│ └── dedup.yaml
└── data
├── demo-dataset.json
└── demo-dataset.jsonl
executor_type 与 ray_address 是切换分布式模式的两个关键配置项,其全局默认值见 data_juicer/config/config_all.yaml:
executor_type: default # executor: "default", "ray", or "ray_partitioned".
ray_address: auto # the address of the Ray cluster.
executor_type 支持 default/local、ray、ray_partitioned 三种取值,由 ExecutorFactory.create_executor 完成调度:设置为 ray 即实例化 RayExecutor;ray_address 默认 auto,可改为集群地址如 ray://<hostname>:<port>。
示例一:运行 Ray 模式数据处理流水线
配置文件 demos/process_on_ray/configs/demo.yaml 将执行器类型设为 ray 并指定自动 Ray 地址,同时声明了 12 个常规 Filter 算子(alphanumeric、average_line_length、character_repetition、flagged_words、language_id_score、maximum_line_length、perplexity、special_characters、stopwords、text_length、words_num、word_repetition):
project_name: 'ray-demo'
dataset_path: './demos/process_on_ray/data/demo-dataset.jsonl' # path to your dataset directory or file
export_path: './outputs/demo/demo-processed'
executor_type: 'ray'
ray_address: 'auto' # change to your ray cluster address, e.g., ray://<hostname>:<port>
# process schedule
# a list of several process operators with their arguments
process:
# Filter ops
- alphanumeric_filter: # filter text with alphabet/numeric ratio out of specific range.
tokenization: false
min_ratio: 0.0
max_ratio: 0.9
- average_line_length_filter: # filter text with the average length of lines out of specific range.
min_len: 10
max_len: 10000
- character_repetition_filter: # filter text with the character repetition ratio out of specific range
rep_len: 10
min_ratio: 0.0
max_ratio: 0.5
- flagged_words_filter: # filter text with the flagged-word ratio larger than a specific max value
lang: en
tokenization: false
max_ratio: 0.0045
flagged_words_dir: ./assets
use_words_aug: false
words_aug_group_sizes: [2]
words_aug_join_char: ""
- language_id_score_filter: # filter text in specific language with language scores larger than a specific max value
lang: en
min_score: 0.8
- maximum_line_length_filter: # filter text with the maximum length of lines out of specific range
min_len: 10
max_len: 10000
- perplexity_filter: # filter text with perplexity score out of specific range
lang: en
max_ppl: 1500
- special_characters_filter: # filter text with special-char ratio out of specific range
min_ratio: 0.0
max_ratio: 0.25
- stopwords_filter: # filter text with stopword ratio smaller than a specific min value
lang: en
tokenization: false
min_ratio: 0.3
stopwords_dir: ./assets
use_words_aug: false
words_aug_group_sizes: [2]
words_aug_join_char: ""
- text_length_filter: # filter text with length out of specific range
min_len: 10
max_len: 10000
- words_num_filter: # filter text with number of words out of specific range
lang: en
tokenization: false
min_num: 10
max_num: 10000
- word_repetition_filter: # filter text with the word repetition ratio out of specific range
lang: en
tokenization: false
rep_len: 10
min_ratio: 0.0
max_ratio: 0.5
运行该示例:
# 从源码运行处理工具
python tools/process_data.py --config demos/process_on_ray/configs/demo.yaml
# 使用命令行工具
dj-process --config demos/process_on_ray/configs/demo.yaml
Data-Juicer 会按配置处理示例数据集,并将结果导出到 export_path 指定的目录。
示例二:运行分布式去重
配置文件 demos/process_on_ray/configs/dedup.yaml 使用 MinHash 去重算子的分布式专用版本 ray_bts_minhash_deduplicator:
project_name: 'demo-dedup'
dataset_path: './demos/process_on_ray/data/'
export_path: './outputs/demo-dedup/demo-ray-bts-dedup-processed'
executor_type: 'ray' # 将执行器类型设置为 "ray"
ray_address: 'auto' # 设置为自动 Ray 地址
# process schedule
# a list of several process operators with their arguments
process:
- ray_bts_minhash_deduplicator: # minhash 去重算子的分布式版本
tokenization: 'character'
运行命令:
# 从源码运行处理工具
python tools/process_data.py --config demos/process_on_ray/configs/dedup.yaml
# 使用命令行工具
dj-process --config demos/process_on_ray/configs/dedup.yaml
该算子除 tokenization 外还支持一批可调参数(默认值见 ray_bts_minhash_deduplicator.py):
| 参数 | 默认值 | 说明 |
|---|---|---|
tokenization |
space |
分词方式:space / punctuation / character / sentencepiece。英文类数据推荐 space,中文类推荐 character,多语言推荐 sentencepiece(需配合 tokenizer_model) |
window_size |
5 | shingling 滑动窗口大小 |
lowercase |
true |
计算 MinHash 前是否转为小写 |
num_permutations |
256 | MinHash 计算中的排列数 |
jaccard_threshold |
0.7 | 近重复判定阈值,两条样本 Jaccard 相似度 >= 该值即视为重复,仅保留其中一条 |
num_bands / num_rows_per_band |
None(自动优化) |
LSH 的 band 数/每 band 行数,默认通过极小化假阳性与假阴性加权和的算法自动确定 |
union_find_parallel_num |
auto |
并查集并行 worker 数,默认取集群 CPU 核数的一半 |
union_threshold |
256 | 触发并查集的 MinHash 分组阈值 |
minhash_batch_size |
auto |
MinHash 计算批大小:CPU 默认 1024,GPU 下按可用显存与 memory_per_sample 自动估算 |
actor_memory |
None |
每个 BTSUnionFind/EdgeBuffer Actor 的内存预留(字节),十亿行规模建议 20GB |
task_memory |
None |
每个 map_batches 任务的内存预留(字节),十亿行规模建议 2GB |
示例三:运行 Ray 分析样例
配置文件 demos/analyze_simple/ray_analyzer.yaml 将 executor_type 设为 ray 以启用分布式 RayAnalyzer:
project_name: 'demo-ray-analyzer'
dataset_path: './demos/data/demo-dataset.jsonl'
executor_type: 'ray' # 将执行器类型设置为 "ray"
ray_address: 'auto' # 设置为自动 Ray 地址
export_path: './outputs/demo-ray-analyzer'
process:
- text_length_filter:
min_len: 10
max_len: 10000
- words_num_filter:
lang: en
min_num: 10
max_num: 10000
- alphanumeric_filter:
min_ratio: 0.25
max_ratio: 0.9
运行该示例:
# 从源码运行分析工具
python tools/analyze_data.py --config demos/analyze_simple/ray_analyzer.yaml
# 使用命令行工具
dj-analyze --config demos/analyze_simple/ray_analyzer.yaml
RayAnalyzer 会计算所有数值统计列的总体统计信息(count、mean、std、min、max)并打印结果;带有统计列的数据集将导出到 export_path 指定的路径,聚合结果另存于工作目录 analysis/overall.csv 与 analysis/overall.md。
总结与实践建议
回顾全文,Data-Juicer 的分布式方案可以概括为“算子引擎无关 + 引擎专属优化”两层结构:日常算子几乎零成本迁移到 Ray;而数据子集分割、JSON 流式读取补丁、基于 BTS 的分布式并查集去重,则针对大规模场景做了专门的性能打磨。结合源码可以进一步得到以下实践要点:
- 小文件多节点场景优先切分数据:当节点数远多于输入文件数时,先运行 tools/data_resplit.py 按 128MB/核数约束切分,能显著降低通信开销;
- 大规模去重优先选用
ray_bts_minhash_deduplicator:它通过分布式并查集与 BTS 负载均衡处理 TB 级等价类合并,单机算子是难以胜任的; - 大规模分析优先用
RayAnalyzer:它用 Ray 原生聚合替代 pandas 物化,可支撑任意规模数据的总体统计;需要图表与相关性分析时,再对采样子集回退到本地 Analyzer; - 多节点运行务必使用共享存储:数据集与导出结果都应放在 NAS 等共享磁盘上,否则各节点无法读到同一份数据;
- 版本兼容注意:JSON 流式读取依赖 PyArrow 的
open_json接口(PyArrow 20.0.0+),旧版本会自动降级为整文件读取,内存表现会有所回退。
延伸阅读
- 分布式处理主文档:docs/Distributed_ZH.md(英文版见 docs/Distributed.md)
- 算子列表与文档:docs/Operators.md、docs/operators/
- Ray 模式执行器源码:ray_executor.py、ray_dataset.py
- 分布式分析源码:ray_analyzer.py、ray_exporter.py
- 分布式去重算子源码:ray_bts_minhash_deduplicator.py
- 数据切分工具:tools/data_resplit.py
- 示例配置与数据:demos/process_on_ray/、demos/analyze_simple/ray_analyzer.yaml