数据处理模板

要在 pipeline 里跑 hadoop/spark/ray/volcano 这类分布式数据处理作业时

05-任务模板 / 模板使用
数据处理hadoopvolcanovolcanojobraysparksparkjob分布式数据处理大数据VC_WORKER_NUMRAY_HOSTSparkApplication

数据处理

以下四类分布式数据处理模板的使用要点。各算子的完整设计与参数见 02-数据处理.md 家族文档,及代码仓库 job-template/job/ 下对应模板目录的 README.md

hadoop 模板

对接客户已有大数据平台(Hadoop / YARN / Hive / Spark)的客户端模板。平台本身不提供大数据集群,只提供可在平台内提交任务的客户端镜像,任务实际提交到你公司已有的集群执行。

  • 使用场景:需要把 spark-submithive -ehadoop fs 等命令提交到公司已有 HDFS / YARN / Hive 集群运行时。
  • 入参--command(要执行的命令,如 spark-submit /path/app.pyhive -e "SELECT 1"hadoop fs -ls /)。
  • 设计:入口 start.sh 先执行 set.sh,用环境变量把集群地址写入 core-site.xml / yarn-site.xml / hive-site.xml,再执行 --command
  • 注意:需按集群实际配置环境变量 FS_DEFAULTFSYARN_RESOURCEMANAGER_ADDRESS(Hive 任务再加 HIVE_METASTORE_URIS);集群组件版本与镜像默认(Hadoop 3.3.6 / Hive 3.1.3 / Spark 3.4.3)不一致时,需按 readme 二开重建客户端镜像。

volcanojob 模板

通过创建 Volcano Job 发起多副本分布式计算任务,所有 Worker 使用相同镜像与命令,由 Volcano 调度器统一调度。

  • 使用场景:数据并行 / 批处理——同一份代码在多机上按序号跑不同分片数据(ETL、离线处理等);或复用已部署 Volcano 集群的队列与 gang 调度能力。
  • 入参--working_dir(工作目录)、--command(由 bash -c 执行)、--num_worker(Worker 数 / Pod 副本数)、--image(Worker 镜像,不填用模板默认)。
  • 设计:launcher 据参数与 KFJ_TASK_* 环境变量组装 Volcano Job(单一 worker task,replicas=num_worker,minAvailable 保证 gang 调度),用 stern 汇聚各 Pod 日志并轮询 Job 状态直至结束。
  • 注意:分布式分片依赖 Volcano 注入的 VC_WORKER_NUM / VC_TASK_INDEX,用户代码据此切分数据(如只处理 index % WORLD_SIZE == RANK);需 kubeflow-pipeline 账号具备相应命名空间下的 CRD / Pod 权限。

ray 模板

基于 Ray 的 Python 多机分布式计算模板,按需在 Kubernetes 上拉起 Ray 集群并执行用户代码。

  • 使用场景:想用 Ray(@ray.remote)做多机并行计算,又不想自己事先部署 Ray 集群时。
  • 入参images(任务镜像,建议基于官方 Ray 镜像封装)、--workdir(工作目录)、--init(可选初始化脚本,Head / Worker 启动前及驱动 Pod 内各执行一次)、--command(用户启动命令,如 python demo.py)、--num_worker(Worker 数,决定集群规模)。
  • 设计:launcher 作为驱动端创建 1 个 Head + N 个 Worker 的 Ray 集群,连接后设置环境变量 RAY_HOST,再在 workdir 执行 --command;用户代码用 <RAY_HOST>:10001 连接同一集群。
  • 注意:镜像需已安装 ray 且与集群版本兼容;用户代码可依据是否存在 RAY_HOST 判断走集群模式还是本地 ray.init() 调试。

sparkjob 模板

通过 Spark Operator 提交 Spark 分布式作业的模板,支持 Java / Python / Scala / R。

  • 使用场景:需要在平台上跑 Spark 分布式计算作业时。
  • 入参--image(执行镜像)、--num_worker(executor 数目)、--code_type(Java / Python / Scala / R)、--code_class(Java / Scala 主类名,其他语言不填)、--code_file(代码文件地址,支持 local:///http:///hdfs:///s3a:///gcs://)、--code_arguments(代码参数)、--sparkConf(每行一个 xx=yy)、--hadoopConf(每行一个 xx=yy)。
  • 设计:把参数组装成 SparkApplication 交给 Spark Operator 调度,按 num_worker 拉起对应数量的 executor 执行代码。
  • 注意:当前版本 Spark Operator 无法挂载分布式存储,code_file 实际只能用 http 地址;分布式存储路径 /mnt/admin/xx/example.py 对应的 http 地址为 http://127.0.0.1/static/mnt/admin/xx/example.py
最后更新 2026-08-12完整文档以官方仓库为准:GitHub Wiki