Data-Juicer 分布式数据处理实战:基于 Ray 的集群化算子执行、子集切分、流式读取与分布式去重

原创2026-10-03 09:58:301,407 阅读
文章标签:人工智能大模型数据工程数据清洗数据增强数据质检

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 监控并导出结果。
  • 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。
  • 例外情况是去重算子:去重算子在单机模式下难以扩展,因此 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)与必要的环境变量。随后:

  1. 由 DatasetBuilder(self.cfg, executor_type="ray") 构建数据集;
  2. 加载 self.cfg.process 中的算子(load_ops),若开启 op_fusion 还会执行算子融合与重排(fuse_operators,相关实现见 op_fusion.py);
  3. 调用 dataset.process(ops, tracer=self.tracer) 执行管线,随后 dataset.data.materialize() 强制物化以收集真实指标(如耗时、输入/输出行数);
  4. 由 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。其工作流为:

  1. 加载数据:通过 Ray 的分布式数据加载(JSON、CSV、Parquet);
  2. 计算统计量:将每个 Filter 的 compute_stats 通过 Ray map_batches 分布式执行——与本地模式完全相同的统计计算逻辑,只是分布式化;
  3. 聚合总体统计:使用 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;
  4. 导出:通过 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),可以进一步理解每项优化背后的实现细节,并据此为自身集群定制切分策略与算子组合。

登录后查看全文
data-juicer