EMR Serverless Custom Job 提供开箱即用的全托管 Job 服务,用户可在 EMR Serverless 提交类似 K8S Job 形态的作业,支持如下功能:
说明
EMR Serverless Custom Job 支持以 SQL 方式提交,可以支持在 DataLeap 中或 EMR Serverless SQL编辑器 中提交作业。这里以 SQL 的方式提交,实际启动一个 Python 作业为例:
打开 EMR Serverless 控制台,并创建任务。
输入下述命令,并点击运行,即可提交一个 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;
Python 代码如下:
print('你好')
在查询日志中,会显示作业实例 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 |
参数名 | 参数含义 | 示例 |
|---|---|---|
serverless.custom.job.task.cores | 每个 Task 的 cpu | 4 |
serverless.custom.job.task.memory | 每个 Task 的 memory | 4Gi |
下表为更新版
参数名 | 参数含义 | 示例 |
|---|---|---|
serverless.tos.volumes |
| [{"type": "fsx", "tosPath": "tos://test-1/tos", "mountPath": "/home/root/", "readOnly": false,"ak":"xxxx", "sk":"xxxxx"}] |
参数名 | 参数含义 | 示例 |
|---|---|---|
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 |
参数名 | 参数含义 | 示例 |
|---|---|---|
las.cluster.type | 节点类型,开启 gpu 需要设置使用 vci 运行 |
|
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 退出后的重启策略,支持如下策略:
默认为 Never | OnFailure |
serverless.custom.job.max.retry | Task 最大重试次数,超过该次数某,标记 Job 失败,默认为3次 |
|
serverless.custom.job.env | 设置 CustomJob 的环境变量,以 JSON array 格式配置。 | ["zzw=123", "wyj=123123123"] |
环境变量名 | 环境变量含义 | 示例 |
|---|---|---|
VC_TASK_INDEX | 当前 Task 索引 | 1 |
VC_TASK_NUM | 当前 Job 的总副本数 | 10 |
EMR Serverless CustomJob支持mount tos路径,您可以将代码托管在tos的某个文件夹下,然后使用mount tos方式来读写代码,无需重打镜像。使用步骤如下:
创建tos路径,并将代码上传。示例中将代码上传至 tos://test-1/tos 路径下:
提交EMR Serverless Custom Job作业,并设置:
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"}];
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需要通过公网下载,请使用挂载辅助网卡功能。
前序准备:
代码编写:
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;
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;
运行后结果成功写入。
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 控制台查看日志。
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。