
人工智能大模型数据工程数据清洗数据增强数据质检【免费下载链接】data-juicerData processing for and with foundation models! ➡️ ➡️ 项目地址https://gitcode.com/gh_mirrors/da/data-juicer点击查看免费下载本指南围绕 Data-Juicer 的数据处理流水线展开系统讲解如何通过命令行dj-process与 Python API 两种方式运行算子流水线并深入剖析算子执行顺序、Op Fusion 性能调优、抽样干跑Dry Run、断点续跑Checkpoint与 Tracing 调试等生产级能力。读完本文你将能够独立编写并运行数据加工 recipe、在训练脚本或 Notebook 中以编程方式调用处理流程、接入 OpenAI 兼容服务与 LiteLLM 完成 API 模型标注并掌握大规模数据集下的调优与故障恢复方法。流水线是如何运行的Data-Juicer 的处理单元是recipe配方一个 YAML 文件其中process字段是一个按顺序排列的算子列表每个算子可以携带自己的参数。运行时的工作方式为读取 recipe解析全局配置与每个算子的参数按列表顺序实例化并执行算子处理输入数据集将结果写入export_path指定的输出文件。整个入口由 tools/process_data.py 承载它根据配置中的executor_type选择不同的执行后端default单机多进程、rayRay 分布式或ray_partitioned分区分发执行对应源码中的DefaultExecutor、RayExecutor与PartitionedRayExecutortools/process_data.py。配置的解析由 data_juicer/config/config.py 中的init_configs完成它基于 jsonargparse 从四个层级合并配置POSIX 风格命令行参数、YAML 配置文件JSON/JSONNet 超集、环境变量、内置默认值。值得注意的细节是解析器只会为 recipe 中实际使用到的算子注册命令行参数_collect_config_info_from_class_docs只针对used_ops生效这让--language_id_score_filter.langen这样的覆盖语法成为可能也让启动速度不受全部 200 算子注册的影响。命令行CLI方式基本用法dj-process --config my-recipe.yamldj-process读取 recipe 后默认按照process列表的先后顺序执行全部算子并把处理结果写到export_path指向的位置。一个最小可运行示例可以参考 demos/process_simple/process.yaml其中同时给出了本地输入与 S3 输入的写法。命令行参数覆盖无需修改 YAML即可在命令行覆盖 recipe 中的任何参数dj-process --config recipe.yaml --np 8 --export_path ./out/result.parquet dj-process --config recipe.yaml --language_id_score_filter.langen第一条命令把并行进程数np调整为 8 并改写输出路径第二条命令只覆盖language_id_score_filter这一个算子的lang参数其余保持 recipe 原样。这套机制在 data_juicer/config/config.py 的update_op_process中实现命令行中出现某算子参数时会先移除 recipe 中对应键再合并保证命令行优先。np的默认值为 4且在初始化时会与机器 CPU 核数比较超限会自动回退data_juicer/config/config.py。自动安装算子依赖dj-install --config my-recipe.yamldj-install会扫描 recipe 中用到的每个算子的源文件通过 AST 解析提取其中的 import 语句与LazyLoader惰性加载声明识别出第三方依赖并预先安装tools/dj_install.py。扫描范围与环境准备细节见 Installation 文档。Python API 方式Python API 比 YAML recipe 提供更细粒度的控制适合训练脚本、Notebook 或自动化流水线。Option 1加载 recipe 运行from data_juicer.config import init_configs from data_juicer.core import DefaultExecutor cfg init_configs(args[--config, my-recipe.yaml]) executor DefaultExecutor(cfg) dataset executor.run()init_configs接收与命令行相同的参数列表返回全局配置对象DefaultExecutor负责加载数据、执行算子并导出结果。其run方法完整流程为preflight 校验 → 加载数据集 → 实例化算子load_ops→ DAG 执行规划 →可选Op Fusion 与自适应 batch size → 逐算子处理 → 导出data_juicer/core/executor/default_executor.py。如果内存中已持有数据集对象例如上游已组装或采样过可以跳过 recipe 的数据源、只复用其算子流水线dataset executor.run(datasetmy_dataset, skip_exportTrue)run的第一个参数传入数据集对象时DatasetBuilder 不再加载数据skip_exportTrue则跳过落盘直接把处理结果返回给调用方。Option 2从配置实例化算子完全不需要 YAML——在 Python 中直接组装算子链from data_juicer.ops import load_ops from data_juicer.core import NestedDataset # Load ops from dict config (same format as YAML process list) ops load_ops([ {language_id_score_filter: {lang: en, min_score: 0.8}}, {text_length_filter: {min_len: 10, max_len: 50000}}, {document_minhash_deduplicator: {tokenization: space, window_size: 5}}, ]) dataset NestedDataset(NestedDataset.from_json(my-data.jsonl)) dataset dataset.process(ops)这里的字典列表与 YAML 中process字段格式完全一致。load_ops会根据算子名从全局注册表OPERATORS.modules中取出对应类并以参数字典实例化同时把原始配置回存到算子的_op_cfg属性data_juicer/ops/load.pyNestedDataset.process(ops)则驱动整个处理过程内部自动完成多进程调度与中间结果缓存。Option 3单算子精细化控制当需要条件逻辑、循环或中间结果检查时可以逐个实例化算子并单独运行from data_juicer.ops.filter import LanguageIDScoreFilter, TextLengthFilter from data_juicer.ops.deduplicator import DocumentMinhashDeduplicator from data_juicer.core import NestedDataset dataset NestedDataset(NestedDataset.from_json(my-data.jsonl)) # Step 1: Language filter lang_filter LanguageIDScoreFilter(langen, min_score0.8) dataset lang_filter.run(datasetdataset) print(fAfter language filter: {len(dataset)} samples) # Step 2: Conditional dedup — only when dataset is large if len(dataset) 10000: dedup DocumentMinhashDeduplicator(tokenizationspace, window_size5) dataset dedup.run(datasetdataset) print(fAfter dedup: {len(dataset)} samples) # Step 3: Length filter length_filter TextLengthFilter(min_len10, max_len50000) dataset length_filter.run(datasetdataset)每个算子类都直接暴露run(dataset...)方法参数与 recipe 中 YAML 键一一对应。以LanguageIDScoreFilter为例它内部使用 FastText 模型识别样本语言与置信度lang可传单个字符串或语言列表空值则只按min_score过滤data_juicer/ops/filter/language_id_score_filter.py。这种模式让先看过滤结果再决定是否去重这样的自适应逻辑变得非常自然。Option 4动态算子组合根据数据特征在运行时决定处理策略适合自动化流水线from data_juicer.ops import load_ops from data_juicer.core import NestedDataset dataset NestedDataset(NestedDataset.from_json(input.jsonl)) # Inspect data to decide processing strategy sample dataset[0] ops_config [] # Add language filter if text field exists if text in sample: ops_config.append({language_id_score_filter: {lang: en, min_score: 0.5}}) # Add image filter if images present if images in sample and sample[images]: ops_config.append({image_shape_filter: {min_width: 256, min_height: 256}}) # Common cleaning ops_config.append({clean_html_mapper: {}}) ops_config.append({text_length_filter: {min_len: 10}}) ops load_ops(ops_config) dataset dataset.process(ops)先通过dataset[0]探查样本字段再按字段存在性动态拼装算子列表最后统一load_ops并process。这种方式同样适用于按批次迭代调整策略的场景。接入 API 模型API ModelsData-Juicer 内置一批依赖托管模型服务的算子典型如extract_keyword_mapper文本关键词提取。使用方式为把api_model设置为服务方提供的模型名并在model_params中配置连接信息。OpenAI 兼容服务首先在已激活的环境中安装 API 依赖uv pip install py-data-juicer[ai_services]然后在 shell 中设置服务地址与密钥export OPENAI_BASE_URLhttps://api.openai.com/v1 export OPENAI_API_KEYyour-api-key使用其他 OpenAI 兼容服务时把 URL 与 Key 换成对应服务的值即可。将下面的 recipe 保存为extract.yamldataset_path指向一个带text字段的 JSONL 文件dataset_path: ./input.jsonl export_path: ./outputs/extracted.jsonl keep_stats_in_res_ds: true process: - extract_keyword_mapper: api_model: gpt-4o-mini model_params: api_backend: openai_compatible sampling_params: temperature: 0运行dj-process --config extract.yaml即可。算子会把抽取出的关键词写入__dj__meta__.keyword字段keep_stats_in_res_ds: true保证该元数据随结果数据集一起导出。这个默认写入键在源码中对应MetaKeys.keyword keyword统一挂在__dj__前缀的 meta 命名空间下data_juicer/utils/constant.py 与 data_juicer/utils/constant.py。openai_compatible是默认后端。也可以在model_params中直接给出base_url和api_key它们的优先级高于同名环境变量。从源码看extract_keyword_mapper的api_model默认值为gpt-4o并支持keyword_key、prompt_template、completion_delimiter、output_pattern、try_num等参数用于自定义输出字段与解析规则data_juicer/ops/mapper/extract_keyword_mapper.py。LiteLLM如果需要通过 LiteLLM 统一路由多家模型把model_params.api_backend设为litellm并使用带供应商前缀的模型名。把上面 recipe 的算子条目替换为process: - extract_keyword_mapper: api_model: openai/gpt-4o-mini model_params: api_backend: litellm sampling_params: temperature: 0随后为所选供应商配置凭证OpenAI 示例使用OPENAI_API_KEY其他供应商遵循各自的 LiteLLM 认证设置当供应商要求时也可以在model_params中直接提供api_key与base_url。选用支持所需任务的算子即可。API 后端支持 chat、embedding 与 Responses 三类端点暴露了api_endpoint参数的算子可以指定走哪个端点。sampling_params用于控制请求选项如 temperaturemodel_params用于配置后端类型、服务地址与凭证。算子执行顺序默认情况下算子自上而下串行执行。顺序对结果和性能都至关重要廉价的过滤器放在前面如text_length_filter、language_id_score_filter——尽早减少样本数量去重放在中间去重依赖全局状态如 MinHash 签名应在初步过滤之后运行昂贵的算子放在最后GPU 推理等重算子只作用于过滤后的子集。process: # Cheap - text_length_filter: { min_len: 10, max_len: 50000 } - language_id_score_filter: { lang: en, min_score: 0.5 } # Dedup - document_minhash_deduplicator: { tokenization: space, window_size: 5 } # Expensive - clean_html_mapper: {} - perplexity_filter: { lang: en, max_ppl: 1500 }注意这个示例中language_id_score_filter的min_score: 0.5比前文示例的 0.8 更低——对于粗筛阶段可以放宽阈值后续如需更高质量再用严格过滤。顺序是否合理还决定了 Op Fusion 的收益下一节会看到 Fusion 会在此基础上进一步重组算子。性能调优Op Fusion算子融合Op Fusion 将兼容的算子融合减少重复的中间计算如反复加载图片、切词等。default与 Ray 执行路径都支持融合吞吐提升幅度取决于 recipe 与数据本身建议针对自己的负载实测op_fusion: true fusion_strategy: probe # probe: group and sort by measured speed; greedy: order by fusion group融合机制在 data_juicer/ops/op_fusion.py 中实现fuse_operators首先按中间变量lines、words、loaded images/audios/videos、sampled frames 六类把共享同一中间变量的 Filter 聚成组组内用FusedFilter合成单算子随后再把连续的 GPU Mapper 融合成单阶段执行data_juicer/ops/op_fusion.py。融合组会改变 recipe 中原有的算子顺序probe策略下default执行器与标准 Analyzer 默认取当前数据集的前 1,000 行数据量更少时取全部行作为探测批次每个算子在该批次的副本上跑一遍每个运行时进程一份副本按实测速度从快到慢重排算子Ray 执行器则按融合组来排布算子。GPU Mapper 融合把连续的 GPU Mapper 融合进单次 GPU 遍历op_fusion: true mapper_fusion: true adaptive_batch_size: truedefault执行器会利用adaptive_batch_size探测并调整批量算子的 batch size。源码中 GPU Mapper 融合还受mapper_fusion_vram_limit默认 0.9约束即一个融合组的聚合显存占用估计不超过单张 GPU 的 90%data_juicer/ops/op_fusion.py。在执行器内部probe策略会先调用adapter.probe_small_batch探测速度再执行fuse_operators融合重排最后用adapt_workloads计算每个算子的自适应 batch sizedata_juicer/core/executor/default_executor.py。抽样干跑Sampling Dry Run对于至少 1,000 行的数据源可以先用default执行器在 1,000 个样本上验证 recipe——方法是不用字符串形式的dataset_path而是改用结构化的dataset配置并保留原来的process列表dataset: max_sample_num: 1000 configs: - type: local path: path/to/your/dataset.jsonl把path换成自己的数据集。max_sample_num设置采样数量小规模试跑时选择不超过源数据量的数值即可。预算更大时会用重复采样填满即样本可重复出现。需要注意采样发生在数据加载之后。如果连输入读取也想省掉就预先切一份小样本文件作为输入。去掉max_sample_num即可处理完整数据集。这套结构化dataset配置同样支持权重混合与远端数据源完整字段说明见 DatasetCfg 文档。Checkpointing 与断点续跑use_checkpoint: truedefault执行器按job_id与工作目录基址定位 checkpoint。首次运行建议指定固定 ID中断后用相同的输入、recipe 与工作目录配置重跑同一条命令即可续跑dj-process --config your-recipe.yaml --use_checkpoint true --job_id recipe-checkpointCheckpoint 存放在解析后的cfg.work_dir/ckpt目录下而work_dir会拼接上job_id。从源码看DefaultExecutor在初始化时即准备CheckpointManager数据集落在ckpt/latest已执行算子记录在ckpt/ckpt_op.json恢复时若记录恰好是当前process列表的前缀就跳过这些已完成算子、从 checkpoint 数据集继续data_juicer/core/executor/default_executor.py 与 data_juicer/utils/ckpt_utils.py。两个重要的互斥/副作用需要留意use_checkpoint会禁用数据缓存HuggingFace datasets 缓存管理因为 checkpoint 本身就是缓存use_checkpoint与op_fusion互斥。源码中一旦开启op_fusion会强制把use_checkpoint置为 Falsedata_juicer/config/config.py配置初始化时还会在use_cacheFalse或use_checkpointTrue的情况下自动关闭 HF 缓存data_juicer/config/config.py。对于ray_partitioned执行器还有更细分的 checkpoint 策略可以指定checkpoint.strategyevery_op/every_n_ops/manual/disabled、checkpoint.n_ops默认每 5 个算子落一次 checkpoint或checkpoint.op_names手动指定恢复时使用--resume job_id它会基于保存的分区边界与内容哈希重建并校验原始分区。完整参数见 GlobalConfig 文档。Tracing 与调试open_tracer: true trace_num: 10开启后每个算子执行前后会对样本进行对比并输出记录帮助定位某个算子到底改了什么、过滤掉了什么。相关可调参数还包括op_list_to_trace只追踪指定算子默认追踪process中全部算子trace_keys指定要包含在 trace 输出中的字段trace_num每个算子抽取展示的样本数量默认 10。Tracer 由DefaultExecutor在初始化时创建trace 结果写入work_dir并支持trace_keys自定义观察字段data_juicer/core/executor/default_executor.py。更多说明见 Tracing 文档。下一步数据分析——在处理前先摸清数据分布数据集配置——输入格式、混合加权、远端数据集全局配置参考——完整参数列表分布式处理——扩展到 Ray 集群算子库——浏览 200 内置算子。赞分享人工智能大模型数据工程数据清洗数据增强数据质检【免费下载链接】data-juicerData processing for and with foundation models! ➡️ ➡️ 项目地址https://gitcode.com/gh_mirrors/da/data-juicer点击查看免费下载相关推荐Data-Juicer数据处理流水线优化操作器融合技术详解Data Juicer数据处理流水线优化操作器融合技术详解 在当今大语言模型LLM快速发展的时代高质量的训练数据已成为决定模型性能的关键因素。Data人工智能大模型数据工程数据清洗数据增强数据质检>data juicer数据处理流水线构建可复用的自动化工作流 引言数据质量的隐形壁垒 在大语言模型LLM训练中数据质量直接决定模型性能上限。然而现实人工智能大模型数据工程数据清洗数据增强数据质检Data-Juicer Agent 流水线最小可运行配置逐项调试指南Data Juicer Agent 流水线最小可运行配置逐项调试指南 导读 Agent 交互数据含 messages、response_choices 的多轮人工智能大模型数据工程数据清洗数据增强数据质检上一篇NemoClaw Inference 模块深度解析模型/Provider 配置、健康探测与 Ollama 认证代理的源码验证机制下一篇chezmoi gopass 模板函数在 dotfiles 中安全注入 gopass 密码凭据创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考