数据处理
以下四类分布式数据处理模板的使用要点。各算子的完整设计与参数见 02-数据处理.md 家族文档,及代码仓库 job-template/job/ 下对应模板目录的 README.md。
hadoop 模板
对接客户已有大数据平台(Hadoop / YARN / Hive / Spark)的客户端模板。平台本身不提供大数据集群,只提供可在平台内提交任务的客户端镜像,任务实际提交到你公司已有的集群执行。
- 使用场景:需要把
spark-submit、hive -e、hadoop fs等命令提交到公司已有 HDFS / YARN / Hive 集群运行时。 - 入参:
--command(要执行的命令,如spark-submit /path/app.py、hive -e "SELECT 1"、hadoop fs -ls /)。 - 设计:入口
start.sh先执行set.sh,用环境变量把集群地址写入core-site.xml/yarn-site.xml/hive-site.xml,再执行--command。 - 注意:需按集群实际配置环境变量
FS_DEFAULTFS、YARN_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。