Data-Juicer 分布式数据处理实战:基于 Ray 的集群化算子执行、子集切分、流式读取与分布式去重
Data-Juicer 分布式数据处理实战:基于 Ray 的集群化算子执行、子集切分、流式读取与分布式去重
Data-Juicer 在单机模式之外提供了完整的 Ray 分布式执行链路:通过 RayDataset 与 RayExecutor 将几乎所有单机算子无缝迁移到 Ray 集群上运行,并针对大规模场景配套了数据子集切分(Subset Splitting)、JSON 流式读取(Streaming Reading)以及基于 MinHash-LSH 与 BTS 算法的分布式去重等引擎级优化。读完本文,你将掌握 Data-Juicer 分布式模式的架构原理、性能优化手段,以及从安装集群到使用 dj-process / dj-analyze 完成分布式处理、去重与数据分析的完整实战流程。
总体概览:基于 Ray 与 PAI 的分布式执行
Data-Juicer 的分布式能力建立在 Ray 之上,并可与阿里云 PAI 等集群环境配合使用。其设计目标是让"单机模式下几乎所有的算子都能无缝地在 Ray 分布式模式下执行",实现这一目标的关键在于将算子的核心处理逻辑与执行引擎解耦:
-
引擎无关的算子实现:绝大多数算子(Filter、Mapper 等)的核心处理函数与执行引擎无关,真正的分布式互操作逻辑集中在两个类中:
- RayDataset:
DJDataset的子类,负责 Ray 数据集(ray.data.Dataset)的加载、预处理与算子调度; - RayExecutor:
BaseExecutor的子类,负责初始化 Ray、构建处理管线、执行 DAG 监控并导出结果。
- RayDataset:
-
Task 与 Actor 双模式支持:在 ray_dataset.py 的
_run_single_op中可以看到,Mapper 和 Filter 会根据算子特性选择执行策略:- 需要维护跨批状态或重资源的算子(如部分去重、LLM 推理算子)使用
ActorPoolStrategy,通过op.__class__配合fn_constructor_args/fn_constructor_kwargs构造 Ray Actor(见 ray_dataset.py 中的_build_actor_pool_strategy,支持固定或弹性并发度区间); - 无状态算子使用
TaskPoolStrategy(size=op.num_proc),以 Task 方式并行执行op.process。
- 需要维护跨批状态或重资源的算子(如部分去重、LLM 推理算子)使用
-
例外情况是去重算子:去重算子在单机模式下难以扩展,因此 Data-Juicer 专门提供了以
ray_xx_deduplicator命名的分布式版本,集中在 data_juicer/ops/deduplicator/ 目录下,包括ray_document_deduplicator、ray_image_deduplicator、ray_video_deduplicator以及基于 BTS 算法的ray_bts_minhash_deduplicator等。
需要说明的是,并非所有算子都天然支持 Ray 模式。在 ray_dataset.py 中,
_run_single_op会检查op._supported_exec_modes是否包含ray,否则抛出NotImplementedError,这是算子在分布式模式下可用的硬性校验条件。
Ray 模式下的数据加载与导出
RayExecutor 初始化时通过 ray_utils.py 中的 initialize_ray 连接集群:默认地址为 auto,并会通过 runtime_env 注入自定义算子路径(custom_operator_paths)与必要的环境变量。随后:
- 由
DatasetBuilder(self.cfg, executor_type="ray")构建数据集; - 加载
self.cfg.process中的算子(load_ops),若开启op_fusion还会执行算子融合与重排(fuse_operators,相关实现见 op_fusion.py); - 调用
dataset.process(ops, tracer=self.tracer)执行管线,随后dataset.data.materialize()强制物化以收集真实指标(如耗时、输入/输出行数); - 由 RayExporter 导出结果,导出路径支持本地磁盘与 S3(
export_path以s3://开头时,会从export_aws_credentials提取凭据注入导出参数,见 ray_executor.py)。
引擎级优化一:数据子集切分(Subset Splitting)
在分布式场景中常会遇到"文件数量少但节点数非常多"的尴尬局面。此时 Ray 默认会根据可用资源将数据集文件切块并分发到所有节点,导致巨大的网络通信开销和 CPU 利用率下降(相关机制可参考 Ray 的 autodetect_parallelism 与输出块调优文档)。
为应对这类场景,Data-Juicer 在读取前会自动将原始数据集预切分为更小的子文件,切分策略兼顾了 Ray 与 Arrow 的读取特性:
- 单个子文件大小目标设定为 128MB(
tools/data_resplit.py中DEFAULT_MAX_FILE_SIZE = 128); - 切分后的子文件数量应保证至少为集群总 CPU 核数的两倍,以便让所有节点、所有核都获得足够粒度的并行任务。
对应的切分工具位于 tools/data_resplit.py,其命令行参数为:
python tools/data_resplit.py \
--ray-address auto \ # Ray 集群地址,默认 auto
--data-dir /path/to/dataset \ # 输入数据集目录或单个 jsonl/json 文件(-i)
--resplit-dir /path/to/output # 切分后子文件输出目录(-o)
从实现看(tools/data_resplit.py),它会先通过 ray.cluster_resources() 获取集群 CPU 总数,再按 max_size = max(1, min(128, total_size / cpu_num / 4))(单位 MB)动态计算单文件大小上限——也就是说当集群核数多、总数据量相对小时,单文件会被切得更小,保证充足的并行度。切分本身也通过 ray.data.from_pandas(df) 配合 map_batches(split_jsonl_dataset) 分布式完成,而不是在单机上串行处理。
遇到大规模节点下性能问题时,可以优先使用该特性;当然也可以按照自己的偏好自行切分数据集。
引擎级优化二:JSON 文件的流式读取
基础模型数据处理中,大量数据集以 JSONL 格式存储且体积巨大,流式读取 JSON 因此成为刚需。然而 Ray Datasets 底层基于 Arrow 的 JSON 读取实现(截至 Ray 2.40、Arrow 18.1.0)并不支持流式读取 JSON 文件。
为弥补这一缺口,Data-Juicer 在 ray_dataset.py 中实现了 JSONStreamDatasource 与 read_json_stream:
- 优先基于 PyArrow 的
pyarrow.json.open_json逐批读取(PyArrow 20.0.0+ 可用),遇到跨块 schema 演化(如某列从 null 变为 string)时会自动unify_schemas统一 schema 并 cast(ray_dataset.py); - 若 PyArrow 版本过旧没有
open_json,则回退到ray.data.read_json以保证兼容性(ray_dataset.py); - 为缓解读取时的 OOM 问题,Data-Juicer 还向上游 Apache Arrow 提交了流式加载补丁(对应仓库 PR #515)。带上该补丁后,Ray 模式默认使用流式加载接口读取 JSON 文件;
- 此外,CSV 与 Parquet 的流式读取已经原生支持。
值得注意的是,read_json_stream 会统一处理 read_options:若用户在配置中以字典形式传入,会转换为 js.ReadOptions 并默认关闭 use_threads(每个 Ray 读取任务本身已是并行的,任务内再开线程会导致 CPU 超订),并接受 json、jsonl、json.gz、jsonl.gz、json.zst、jsonl.zst 等常见扩展名(ray_dataset.py)。RayDataset.read 在数据格式为 json/jsonl 系列时即路由到该流式读取接口。
引擎级优化三:分布式去重(MinHash-LSH + BTS)
单机去重难以扩展到海量数据,Data-Juicer 因此在 Ray 模式下提供了一族优化过的分布式去重算子(data_juicer/ops/deduplicator/),核心思路包括:
- 多进程 Union-Find 集合:在 Ray Actor 中实现多进程并查集,用于跨分区维护去重状态;
- 负载均衡的 BTS 算法:采用 BTS(一种分布式等价类合并算法)完成等价类合并,避免单一节点成为瓶颈。
去重算子在 ray_dataset.py 中走独立分支:Deduplicator 与 Pipeline 类型的算子直接调用 op.run(self.data),由算子自身管理分布式执行逻辑。
官方实验表明,该算子在 1280 个 CPU 核上可在 3 小时内完成 TB 级数据集的去重;与此类去重算子的朴素版本相比,针对 Ray 模式的专项优化可获得 2x~3x 的加速。
性能表现
变规模数据处理实验
官方在包含数十亿样本的数据集上进行了实验:以 56 万样本的多模态数据集为基准,按 1x~125000x 的倍数扩充构造不同规模的数据集。结果显示 Data-Juicer 的 Ray 模式表现出良好的可扩展性(Scalability),即处理耗时随数据规模近似线性增长。
大规模数据集上的分布式去重
官方使用基于 MinHash 的 RayDeduplicator 在 200GB、1TB、5TB 三个数据规模、640~1280 核范围内进行了测试:
| # CPU | 200GB Time | 1TB Time | 5TB Time |
|---|---|---|---|
| 4 * 160(640 核) | 11.13 min | 50.83 min | 285.43 min |
| 8 * 160(1280 核) | 7.47 min | 30.08 min | 168.10 min |
从表中可以读出两个规律:
- 数据量 5x 增长时,处理耗时增长约 4.02x~5.62x(接近线性,说明该算子在大数据量下扩展性良好);
- CPU 核数翻倍时,处理耗时降至原来的 58.9%~67.1%(接近理想的 50% 加速比)。
注:以上实验数据均来自 docs/Distributed.md 的官方记录,属于特定软硬件环境下的实测结果,实际性能会随集群配置、数据分布与算子组合而变化。
分布式数据分析:RayAnalyzer
除了分布式数据处理,Data-Juicer 还通过 RayAnalyzer 提供分布式数据分析能力,它是本地 Analyzer 的分布式对应物。
工作原理
RayAnalyzer 复用本地 Analyzer 的 dj-analyze CLI 入口:当配置中 executor_type 设为 ray 时,dj-analyze 自动分发到 RayAnalyzer。其工作流为:
- 加载数据:通过 Ray 的分布式数据加载(JSON、CSV、Parquet);
- 计算统计量:将每个 Filter 的
compute_stats通过 Raymap_batches分布式执行——与本地模式完全相同的统计计算逻辑,只是分布式化; - 聚合总体统计:使用 Ray 原生聚合算子(
Mean、Std、Min、Max,以及Count)计算总体统计,不物化到 pandas,因此适用于任意规模的数据集。从实现看(ray_analyzer.py),聚合前会先把 stats/meta 字典列通过_flatten_stats_batch展开为独立列,仅对数值型列(含数值列表列,列表列先用pc.list_flatten展开)计算 count/mean/std/min/max,非数值列跳过;结果会同时输出到<work_dir>/analysis/overall.csv与overall.md; - 导出:通过 RayExporter 将带统计列的数据集导出到磁盘。
另外,RayAnalyzer 支持 auto 模式:当配置开启 auto 时,会按 auto_num 对数据进行采样后再分析(ray_analyzer.py),适合在超大数据集上做快速概览。
与本地 Analyzer 的差异
| 功能 | Local Analyzer | RayAnalyzer |
|---|---|---|
| 总体统计(count/mean/std/min/max) | 支持 | 支持 |
| 逐列分布图 | 支持 | 不支持 |
| 相关性分析 | 支持 | 不支持 |
| 百分位数 | 支持 | 不支持 |
| 可扩展性 | 单机 | 分布式(Ray 集群) |
因此官方建议:在规模化场景下用 RayAnalyzer 计算总体统计;需要细粒度可视化与相关性分析时,可在采样子集上使用本地 Analyzer。
快速开始
第一步:安装与启动 Ray 集群
先安装 Data-Juicer 及其分布式依赖:
# 安装 Data-Juicer 最小依赖
uv pip install -v -e .
# 安装包含 Ray 及其他分布式库的依赖
uv pip install -v -e ".[distributed]"
然后启动 Ray 集群(更详细的启动方式可参考 Ray 官方文档的 starting-ray 章节):
# 以 head 节点方式启动集群
ray start --head
# (可选)在其他机器上连接到集群
ray start --address='{head_ip}:6379'
仓库在 demos/process_on_ray/ 目录下提供了两个配置文件与两个测试数据集:
demos/process_on_ray
├── configs
│ ├── demo.yaml
│ └── dedup.yaml
└── data
├── demo-dataset.json
└── demo-dataset.jsonl
[!Important] 如果在多节点上运行这些 demo,需要将 demo 数据集放到共享磁盘(如 NAS)上,并同步修改配置中的
dataset_path与export_path,让所有节点都能访问输入并写入共享的输出目录。
第二步:Ray 模式数据处理示例
在 demo.yaml 中,executor_type 被设为 ray,ray_address 设为自动探测:
project_name: 'ray-demo'
dataset_path: './demos/process_on_ray/data/demo-dataset.jsonl' # 数据集目录或文件路径
export_path: './outputs/demo/demo-processed'
executor_type: 'ray' # 将执行器类型设为 "ray"
ray_address: 'auto' # 自动获取 Ray 地址;也可改为 ray://<hostname>:<port>
# process schedule
# 一组处理算子及其参数
process:
- alphanumeric_filter: # 过滤字母/数字比例超出范围的文本
tokenization: false
min_ratio: 0.0
max_ratio: 0.9
- average_line_length_filter: # 过滤平均行长超出范围的文本
min_len: 10
max_len: 10000
- character_repetition_filter: # 过滤字符级重复率超出范围的文本
rep_len: 10
min_ratio: 0.0
max_ratio: 0.5
- flagged_words_filter: # 过滤敏感词比例过大的文本
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: # 按语言识别分数过滤文本
lang: en
min_score: 0.8
- maximum_line_length_filter: # 过滤最大行长超出范围的文本
min_len: 10
max_len: 10000
- perplexity_filter: # 按困惑度分数过滤文本
lang: en
max_ppl: 1500
- special_characters_filter: # 过滤特殊字符比例超出范围的文本
min_ratio: 0.0
max_ratio: 0.25
- stopwords_filter: # 过滤停用词比例过小的文本
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: # 过滤文本长度超出范围的样本
min_len: 10
max_len: 10000
- words_num_filter: # 过滤词数超出范围的样本
lang: en
tokenization: false
min_num: 10
max_num: 10000
- word_repetition_filter: # 过滤词级重复率超出范围的文本
lang: en
tokenization: false
rep_len: 10
min_ratio: 0.0
max_ratio: 0.5
该配置共包含 12 个常规算子。关于各算子的完整参数含义与取值范围,可对照 docs/operators/filter/ 下的算子文档进一步查阅。
运行该 demo:
# 方式一:从源码运行工具
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 会按配置处理 demo 数据集,并将结果导出到 export_path 指定的目录。
第三步:分布式去重示例
在 dedup.yaml 中,执行器同样设为 ray,并启用了专门的分布式 MinHash 去重算子:
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
# 一组处理算子及其参数
process:
- ray_bts_minhash_deduplicator: # 分布式版本的 minhash 去重器
tokenization: 'character'
运行该 demo:
# 方式一:从源码运行工具
python tools/process_data.py --config demos/process_on_ray/configs/dedup.yaml
# 方式二:使用命令行工具
dj-process --config demos/process_on_ray/configs/dedup.yaml
提示:
tokenization: 'character'表示按字符粒度进行 MinHash 分词,适合中文等无空格分隔的语言;若处理英文等空格分隔文本,可改用'word'。ray_bts_minhash_deduplicator对应实现为 ray_bts_minhash_deduplicator.py,仓库中另有基于 C++ 加速的 ray_bts_minhash_cpp_deduplicator.py 可供选用。
第四步: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
运行该 demo:
# 方式一:从源码运行工具
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 指定的路径;总体统计结果还会以 overall.csv / overall.md 形式写入 work_dir/analysis/ 目录。
小结
Data-Juicer 的分布式能力是一套"单机算子 + 分布式执行层"的完整方案:绝大多数算子无需重写即可通过 RayDataset / RayExecutor 在 Ray 集群上运行;针对大规模场景,官方还提供了子集切分(保证子文件数与 CPU 核数的比例关系)、JSON 流式读取(规避 OOM)与 BTS 分布式去重三项引擎级优化。配合 dj-process / dj-analyze 命令行入口和 demos/process_on_ray 下的开箱即用配置,用户可以快速从单机实验平滑过渡到百亿级样本的云端分布式处理。深入阅读 docs/Distributed.md 及相关源码(ray_dataset.py、ray_executor.py、ray_analyzer.py、tools/data_resplit.py),可以进一步理解每项优化背后的实现细节,并据此为自身集群定制切分策略与算子组合。