You need to enable JavaScript to run this app.
文档中心
文档控制台
注册
E-MapReduce

E-MapReduce

复制全文
下载 pdf
提交作业
EMR Serverless Custom Job 使用说明
复制全文
下载 pdf
EMR Serverless Custom Job 使用说明

EMR Serverless Custom Job 提供开箱即用的全托管 Job 服务,用户可在 EMR Serverless 提交类似 K8S Job 形态的作业,支持如下功能:

  • 无需考虑底层集群运维,开箱即用,可选择 按量付费 或 包年包月。
  • 支持自定义镜像,可以使用您镜像仓库下的镜像。
  • 支持自定义 Job 副本数,重试策略等参数。
  • 支持挂载辅助网卡,可以与您指定 VPC 下的资源网络互通。
  • 支持 Mount TOS 路径至 Job 中,您可以在 Job 中使用访问本地路径的方式无缝访问 TOS。

前提条件
  1. 开通 EMR Serverless Spark 服务。
  2. 上传自定义镜像至您账户下的镜像仓库。
  3. 确认镜像仓库功能已开白,如未开白请联系火山工作人员。

任务提交

说明

  • Job:一个 Job 对应 EMR Serverless 的一个任务,Job 为 Task 的集合。
  • Task:Job 中,每个副本都称为 Task。

EMR Serverless Custom Job 支持以 SQL 方式提交,可以支持在 DataLeap 中或 EMR Serverless SQL编辑器 中提交作业。这里以 SQL 的方式提交,实际启动一个 Python 作业为例:

  1. 打开 EMR Serverless 控制台,并创建任务。

  2. 输入下述命令,并点击运行,即可提交一个 EMR Serverless CustomJob。

    说明

    下述 case 可以提交一个自定义的 Python 任务,并使用3副本运行。

    set serverless.custom.job.image = <您的镜像地址>;
    set serverless.custom.job.entrypoint = ["python", "<script_path>"];
    set serverless.custom.job.replicas = 3;
    
  3. Python 代码如下:

print('你好')
  1. 在查询日志中,会显示作业实例 id 和作业状态。

    您可以到上级目录的作业实例列表,根据实例id查询作业情况,并查看日志信息(当前日志仅显示多副本中索引为0的 Task 日志)。

详细参数列表

核心提交参数

参数名

参数含义

示例

serverless.custom.job.image

Job 镜像地址,填写您镜像仓库下的镜像地址

image-repo-cn-beijing.cr.volces.com/emr-serverless/job:v1.0

serverless.custom.job.entrypoint

Job 启动命令,需要用户以 JSON array 格式配置

["python", "tos/redis_test.py"]或["sh", "-c", "python /home/root/a.py"]

serverless.custom.job.args

Job 启动参数,需要用户以 JSON array 格式配置

["--url_list", "/home/root/img2dataset/sbu-captions-all.json", "--out_folder", "/home/root/img2dataset/result/sbu-captions"]

serverless.custom.job.replicas

Job 副本数,默认为1

3

自定义Task资源使用

参数名

参数含义

示例

serverless.custom.job.task.cores

每个 Task 的 cpu

4

serverless.custom.job.task.memory

每个 Task 的 memory

4Gi

Mount TOS

下表为更新版

参数名

参数含义

示例

serverless.tos.volumes

  • 指定 TOS 挂载配置,支持 FSX\S3FS两种挂载方式,需配置:
    • type:挂载方式,

      说明

      type,支持设置为 fsx 或 s3fs。

    • tosPath:TOS路径,
    • mountPath:挂载路径,
    • readOnly:是否只读,
    • ak/sk:访问TOS的AK\SK。
  • 需要用户以 JSON array 格式配置。
  • 详细参数信息请参见:挂载 TOS

[{"type": "fsx", "tosPath": "tos://test-1/tos", "mountPath": "/home/root/", "readOnly": false,"ak":"xxxx", "sk":"xxxxx"}]

跨VPC挂载辅助网卡,访问用户VPC下资源

参数名

参数含义

示例

serverless.cross.vpc.access.enabled

是否开启跨VPC

true

serverless.cross.vpc.vpc.id

访问资源所在的VPC id

vpc-xxx

serverless.cross.vpc.subnet.ids

访问vpc下资源所用的子网id,支持多选,逗号分隔

subnet-aa,subnet-bb

serverless.cross.vpc.security.group.id

访问vpc下资源所用的安全组id

sg-xxxx

GPU支持

参数名

参数含义

示例

las.cluster.type

节点类型,开启 gpu 需要设置使用 vci 运行

  • vci
  • vke

serverless.custom.job.gpu.enabled

开启 gpu 支持

true

serverless.custom.job.task.gpus

每个 Task 的 GPU卡 数

1

其他

参数名

参数含义

示例

serverless.custom.job.working.dir

Task 默认的工作路径,默认为 /home/root

/home/root/test

serverless.custom.job.restart.policy

Task 退出后的重启策略,支持如下策略:

  • Always:一直重启,即使Task正常退出
  • OnFailure:失败才重启
  • Never:从不重启

默认为 Never

OnFailure

serverless.custom.job.max.retry

Task 最大重试次数,超过该次数某,标记 Job 失败,默认为3次

  • 3:Task重试3次后仍然失败,则标记Job失败
  • -1:Task无限重试

serverless.custom.job.env

设置 CustomJob 的环境变量,以 JSON array 格式配置。

["zzw=123", "wyj=123123123"]

支持的环境变量

环境变量名

环境变量含义

示例

VC_TASK_INDEX

当前 Task 索引

1

VC_TASK_NUM

当前 Job 的总副本数

10

最佳实践

使用TOS托管代码

EMR Serverless CustomJob支持mount tos路径,您可以将代码托管在tos的某个文件夹下,然后使用mount tos方式来读写代码,无需重打镜像。使用步骤如下:

  • 创建tos路径,并将代码上传。示例中将代码上传至 tos://test-1/tos 路径下:

  • 提交EMR Serverless Custom Job作业,并设置:

    • serverless.tos.volumes,将 tos 文件夹挂载到目标 dir 下
    • serverless.custom.job.entrypoint=["python", "tos/redis_test.py"],设置启动命令为启动 Python 脚本,脚本路径为 tos/redis_test.py ,对应 tos://test-1/tos/redis_test.py 文件
      完整代码如下:
    set serverless.custom.job.image=emr-vke-qa-cn-beijing.cr.volces.com/emr/custom-job:python3.10-redis;
    set serverless.custom.job.entrypoint=["python", "tos/redis_test.py"];
    set serverless.emr.access.key=xx;
    set serverless.emr.secret.key=xxx;
    set serverless.tos.volumes=[{"type": "fsx", "tosPath": "tos://test-1/tos", "mountPath": "/home/root/", "readOnly": false,"ak":"xxxx", "sk":"xxxxx"}];
    

使用多副本,加速img2dataset下载

EMR Serverless CustomJob 支持多副本,并在环境变量中可获取当前副本数,以及Task所在的副本索引
img2dataset 是一个批量下载图片的工具,在实际应用中,您可以将包含 URL 的数据文件手动分片,并配合 Custom Job 的多副本机制并行执行,以实现加速下载。步骤如下:
以官方case https://github.com/rom1504/img2dataset/blob/main/dataset_examples/SBUcaptions.md ,处理 SBU Captions 数据集中的一小部分 URL 为例:

说明

EMR Serverless CustomJob 默认无法访问公网,如果URL需要通过公网下载,请使用挂载辅助网卡功能。

  • 前序准备:

    • 将要下载的 URL json 列表放在 TOS 上,这里路径为 tos://test-1/img2dataset/sbu-captions-all.json。
  • 代码编写:

    • 读取 URL 列表,并对 URL 进行 hash 取余,判断 url是 否可以被当前 Task 处理。
    • 通过 img2dataset 下载,并写入本地路径。
    import argparse
    import json
    import logging
    import os
    
    import img2dataset
    from fsspec.implementations.local import LocalFileSystem
    
    
    def download(url_list, output_folder):
        logging.info(f"start task, input url dir: {url_list}, output dir: {output_folder}")
        img2dataset.download(
            processes_count=16,
            thread_count=64,
            retries=0,
            timeout=10,
            url_list=url_list,
            image_size=256,
            resize_only_if_bigger=True,
            output_folder=output_folder,
            output_format="webdataset",
            input_format="json",
            url_col="image_urls"
        )
    
    
    if __name__ == "__main__":
        parser = argparse.ArgumentParser()
        parser.add_argument("--url_list", default="/home/root/img2dataset/sbu-captions-all.json")
        parser.add_argument("--out_folder", default="/home/root/img2dataset/result/sbu-captions")
        args = parser.parse_args()
        task_index = os.getenv("VC_TASK_INDEX", 1)
        task_num = os.getenv("VC_TASK_NUM", 1)
        logging.info(f"current task num: {task_num}, index: {task_index}")
    
        task_list = []
        # assign task
        fs = LocalFileSystem()
        with fs.open(args.url_list, 'r') as f:
            local_dict = json.loads(f.readline())
        for url in local_dict['image_urls']:
            if hash(url) % int(task_num) == int(task_index):
                task_list.append(url)
    
        # start task
        logging.info(f"index: {task_index}, len url: {len(task_list)}")
        url_path = f'url-{task_index}.json'
        with fs.open(url_path, 'w') as f:
            json_dict = {'image_urls': task_list}
            f.write(json.dumps(json_dict))
            f.close()
    
        download(url_path, args.out_folder)
    
  • 提交作业,完整代码如下:

    set serverless.custom.job.image=emr-vke-qa-cn-beijing.cr.volces.com/emr/custom-job:python3.10-img;
    set serverless.custom.job.entrypoint=["python", "tos/image2dataset_test.py", "--url_list", "/home/root/img2dataset/sbu-captions-all.json", "--out_folder", "/home/root/img2dataset/result/sbu-captions"];
    set serverless.emr.access.key=xx;
    set serverless.emr.secret.key=xx;
    set serverless.tos.volumes=[
     {"type": "fsx", "tosPath": "tos://test-1/tos", "mountPath": "/home/root/tos", "readOnly": false,"ak":"xxxx", "sk":"xxxxx"},
     {"type": "fsx", "tosPath": "tos://test-1/img2dataset", "mountPath": "/home/root/img2dataset", "readOnly": false,"ak":"xxxx", "sk":"xxxxx"}
    ];
    set serverless.custom.job.replicas = 2;
    

读写租户VPC下的redis

EMR ServerlessCustom Job 运行的网络环境为EMR 托管环境,默认与您账户下的 VPC 网络不互通。如果希望网络能打通,需要设置参数,通过挂载辅助网卡方式进行网络打通。
以读写您VPC下的全托管redis为例:

  • 准备工作:

    • 开通 redis

    • 编写 Python 代码

      import redis
      
      # 建立到 Redis 服务器的连接
      r = redis.Redis(host='redis-cnlf68vdfvox694ma.redis.ivolces.com', port=6379, db=0, password='11111111')
      
      # 写入数据
      r.set('my_key', 'Hello, Redis!')
      
      # 读取数据
      value = r.get('my_key')
      print(value.decode('utf-8'))  # 打印读取的值
      
      # 也可以使用哈希表进行读写操作
      r.hset('my_hash', 'field1', 'value1')
      print(r.hget('my_hash', 'field1').decode('utf-8'))
      
      # 关闭连接
      r.close()
      
    • 在作业中额外设置vpc相关参数,需要设置您的vpc id、子网id(支持多子网)、vpc下的安全组id

      set serverless.cross.vpc.access.enabled = true;
      set serverless.cross.vpc.vpc.id = vpc-rs09o6sxx;
      set serverless.cross.vpc.subnet.ids = subnet-13fxx,subnet-14fxx;
      set serverless.cross.vpc.security.group.id = sg-rsxx;
      
  • 运行后结果成功写入。

使用Dataleap提交CustomJob

EMR Serverless CustomJob 支持使用 Dataleap 提交,您可以在 Dataleap 上创建一个 EMR Serverless Spark SQL 任务。

并通过SQL方式提交作业:
这里原来是一张图,由于需要更新参数所以改成代码块,copy 了原来图上的内容改了下参数

set serverless.custom.job.image=emr-vke-qa-cn-beijing.cr.volces.com/emr/custom-job:python3.10-img;
set serverless.emr.access.key=xx;
set serverless.emr.secret.key=xx;
set serverless.custom.job.entrypoint=["python","tos/custom_job.py"];
set serverless.tos.volumes=[{"type": "fsx", "tosPath": "tos://test-1/tos", "mountPath": "/home/root/", "readOnly": false,"ak":"xxxx", "sk":"xxxxx"}];

set serverless.cross.vpc.access.enabled=true; 
set serverless.cross.vpc.id=<vpc_id>; 
set serverless.cross.vpc.subnet.ids=subnet-mimeqvc0gtmo5smt1ayqtfwh; 
set serverless.cross.vpc.security.group.id=sg-miy0rjg55am85smt1bvecy88;

在调试日志中,可以找到 EMR Serverless 的作业实例 id,并通过实例 id 在 EMR Serverless 控制台查看日志。

使用GPU

EMR Serverless CustomJob 支持使用按量付费的 GPU,您可以通过参数打开 GPU 的支持。
完整 case 如下:

set serverless.custom.job.image=emr-vke-qa-cn-beijing.cr.volces.com/emr/custom-job:python312;
set serverless.custom.job.entrypoint=["python", "tos/custom_job.py"];
set serverless.emr.access.key=xx;
set serverless.emr.secret.key=xx;
set serverless.tos.volumes=[{"type": "fsx", "tosPath": "tos://test-1/tos", "mountPath": "/home/root/", "readOnly": false,"ak":"xxxx", "sk":"xxxxx"}];

set las.cluster.type = vke;
set serverless.custom.job.gpu.enabled = true;
set serverless.custom.job.node.selector = gpu-pool;

运行后,可以看到 Job 默认挂载了 GPU。

最近更新时间:2026.07.06 17:44:17
这个页面对您有帮助吗?
有用
有用
无用
无用