Ray集群是EMR Serverless Spark工作空间提供的分布式计算框架,支持Python原生分布式计算、机器学习模型训练与推理等场景。本文介绍如何创建、启动Ray集群,以及如何提交Ray作业。
创建Ray集群
登录EMR Serverless Spark控制台,进入目标工作空间。
在左侧导航栏,单击集群管理,然后单击RAY集群页签。
单击创建Ray集群,在创建面板中配置以下参数,然后单击创建。
集群名称:输入集群名称。
引擎版本:选择引擎版本。当前提供默认版本err-1.0.1 (Ray 2.47.1, Python 3.12)。
节点组:如需同时使用 CPU 与 GPU 混合资源,可分别创建 CPU 节点组和 GPU 节点组。单击+添加Worker节点组,配置以下信息。
节点组名称:输入节点组名称,不同节点组的名称不能重复。
资源队列:选择资源队列。
资源规格:选择节点规格。
节点数量:设置Worker节点数量。
(可选)网络连接:如需从Ray集群访问同VPC内的服务,选择已创建的网络连接。
(可选)挂载纳管文件目录:选择已创建的纳管文件目录,以实现在 Ray 任务中读写文件目录下的文件。挂载是一种常用的数据访问方式,支持将 CPFS、NAS 和 OSS 路径映射为本地目录,方便程序以本地路径的方式访问数据。如需创建新的纳管文件目录,详情请参见纳管文件目录。
(可选)集群高级配置:以JSON格式配置高级参数,详情请参见集群高级配置参数说明。
当前控制台暂不支持配置Worker节点自动伸缩,如需配置,可通过API方式修改,详情请参见CreateRayCluster - 创建Ray集群。
启动Ray集群
在Ray集群列表中,找到目标集群,单击启动。
等待集群运行状态变为运行中。
说明工作空间内首次启动Ray集群时,系统需要创建额外组件,启动时间约需2~3分钟;后续启动速度会明显加快。
集群启动后,单击调用信息,记录提交Ray作业所需的地址和Token信息。
警告当前版本中,Token对同一Ray集群永久有效,请妥善保管,避免泄露。
单击Dashboard,可进入Ray集群的监控界面,查看集群资源使用情况。
提交Ray作业
方式一:交互式开发
适合快速入门和调试。在EMR Serverless Spark控制台的Notebook环境中直接连接Ray集群(无需额外配置网络)。如从本地连接,需确保本地机器与Ray集群之间网络互通,且环境为Python 3.12。
安装Ray客户端:
pip install ray[client]==2.47.1在集群的调用信息页面,获取gRPC 地址和Token。
使用以下示例代码连接Ray集群,进行交互式开发:
import ray import os def get_metadata(): headers = {"ray-token": "<yourToken>"} return [(key.lower(), value) for key, value in headers.items()] ray.init(address="<your_gRPC_address>", _metadata=get_metadata()) import time @ray.remote def square(x): time.sleep(0.1) return x * x futures = [square.remote(i) for i in range(10)] results = ray.get(futures) print("平方结果:", results)
方式二:SDK提交(批任务)
适合将Ray作业嵌入工程代码中提交。本地环境要求:Python 3.12。
安装Ray Job Submission客户端库:
pip install "ray[default]==2.47.1"在集群的调用信息页面,获取公网调用地址和Token。
使用以下示例代码连接Ray集群并提交作业:
import time from ray.job_submission import JobSubmissionClient custom_headers = { "ray-token": "<yourToken>" } client = JobSubmissionClient("<公网调用地址>", headers=custom_headers) # 或者用内网,需要保证提交端在同一region vpc下,建议内网提交,更加稳定 # client = JobSubmissionClient("http://emr-spark-ray-gateway-cn-beijing-internal.spark.emr.aliyuncs.com", headers=custom_headers) job_id = client.submit_job( entrypoint="python -c 'print(\"Hello from Ray Client!\")'" ) print(f"Submitted job with ID: {job_id}") while True: status = client.get_job_status(job_id) print(f"Job status: {status}") if status.is_terminal(): break time.sleep(1) logs = client.get_job_logs(job_id) print("Job logs:") print(logs)
(可选)如需在作业中访问挂载目录中的文件,可在创建Ray集群时挂载纳管文件目录,然后在提交作业时通过挂载路径引用文件。以下示例假设已将OSS目录挂载到
/mnt/myoss路径下:import time from ray.job_submission import JobSubmissionClient custom_headers = { "ray-token": "<yourToken>" } client = JobSubmissionClient("<公网调用地址>", headers=custom_headers) # 通过挂载路径直接引用OSS中的脚本文件作为entrypoint job_id = client.submit_job( entrypoint="python /mnt/myoss/main.py" ) print(f"Submitted job with ID: {job_id}") while True: status = client.get_job_status(job_id) print(f"Job status: {status}") if status.is_terminal(): break time.sleep(1) logs = client.get_job_logs(job_id) print("Job logs:") print(logs)说明在Ray任务代码中,也可以直接通过挂载路径读写文件。例如,读取
/mnt/myoss/test.txt或写入文件到/mnt/myoss/目录下,写入的文件会同步到对应的OSS路径中。以下示例展示如何在Ray任务中读写挂载目录中的文件:
# 下面是提交到ray集群的代码 import ray import time ray.init() @ray.remote def read_from_mount_path(i): with open('/mnt/myoss/test.txt', 'r', encoding='utf-8') as f: content = f.read() return content @ray.remote def write_to_mount_path(i): content = "Hello world " + str(i) with open('/mnt/myoss/output' + str(i) + '.txt', 'w', encoding='utf-8') as f: f.write(content) return '写入成功' print("--- 开始运行远程任务 ---") start_time = time.time() obj_refs = [read_from_mount_path.remote(i) for i in range(4)] obj_refs2 = [write_to_mount_path.remote(i) for i in range(4)] results = ray.get(obj_refs) results2 = ray.get(obj_refs2)
方式三:命令行提交(批任务)
适合与外部系统集成。本地环境要求:Python 3.12。
安装Ray Job Submission客户端库:
pip install "ray[default]==2.47.1"在集群的调用信息页面,获取公网调用地址和Token。
执行以下命令提交Ray作业:
# ray job submit前必须设置ray address和headers export RAY_ADDRESS='<公网调用地址>' export RAY_JOB_HEADERS='{"ray-token": "<yourToken>"}' ## 提交任务 ### working dir 支持本地 & Ray core distributed sort ray job submit --working-dir "." -- python test-ray-core-sort.py ### 支持OSS读写 ray job submit --working-dir "." -- python test-saving-data-test.py ### 支持OSS-HDFS读写 ray job submit --working-dir "." -- python test-saving-data-oss-hdfs-test.py ### 代码支持Mount读写 ray job submit --working-dir "." -- python test-saving-data-test-mount.py
更多命令行参数,请参见Ray Job Submission CLI官方文档。
集群高级配置参数说明
集群高级配置以JSON格式填写,所有参数均为可选。
参数 | 是否必选 | 说明 |
| 可选 | 指定在集群启动时需要下载到Head节点和Worker节点的OSS文件。支持OSS和OSS-HDFS路径,多个路径之间用英文逗号(,)分隔。文件下载至各节点的 |
| 可选 | 指定用于初始化Head节点和Worker节点Python基础环境的requirements.txt文件。支持OSS和OSS-HDFS路径,路径格式必须为 |