全部产品
Search
文档中心

开源大数据平台E-MapReduce:向Ray集群提交任务

更新时间:Jun 09, 2026

Ray集群是EMR Serverless Spark工作空间提供的分布式计算框架,支持Python原生分布式计算、机器学习模型训练与推理等场景。本文介绍如何创建、启动Ray集群,以及如何提交Ray作业。

创建Ray集群

  1. 登录EMR Serverless Spark控制台,进入目标工作空间。

  2. 在左侧导航栏,单击集群管理,然后单击RAY集群页签。

  3. 单击创建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集群

  1. 在Ray集群列表中,找到目标集群,单击启动

  2. 等待集群运行状态变为运行中

    说明

    工作空间内首次启动Ray集群时,系统需要创建额外组件,启动时间约需2~3分钟;后续启动速度会明显加快。

  3. 集群启动后,单击调用信息,记录提交Ray作业所需的地址和Token信息。

    警告

    当前版本中,Token对同一Ray集群永久有效,请妥善保管,避免泄露。

    单击Dashboard,可进入Ray集群的监控界面,查看集群资源使用情况。

提交Ray作业

方式一:交互式开发

适合快速入门和调试。在EMR Serverless Spark控制台的Notebook环境中直接连接Ray集群(无需额外配置网络)。如从本地连接,需确保本地机器与Ray集群之间网络互通,且环境为Python 3.12。

  1. 安装Ray客户端:

    pip install ray[client]==2.47.1
  2. 在集群的调用信息页面,获取gRPC 地址Token

  3. 使用以下示例代码连接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)

    image

方式二:SDK提交(批任务)

适合将Ray作业嵌入工程代码中提交。本地环境要求:Python 3.12。

  1. 安装Ray Job Submission客户端库:

    pip install "ray[default]==2.47.1"
  2. 在集群的调用信息页面,获取公网调用地址Token

  3. 使用以下示例代码连接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) 

    image

  4. (可选)如需在作业中访问挂载目录中的文件,可在创建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。

  1. 安装Ray Job Submission客户端库:

    pip install "ray[default]==2.47.1"
  2. 在集群的调用信息页面,获取公网调用地址Token

  3. 执行以下命令提交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格式填写,所有参数均为可选。

参数

是否必选

说明

userDefinedFiles

可选

指定在集群启动时需要下载到Head节点和Worker节点的OSS文件。支持OSS和OSS-HDFS路径,多个路径之间用英文逗号(,)分隔。文件下载至各节点的/home/ray/work-dir目录下。示例:oss://mybucket/hello.py,oss://mybucket2/test/test.jar

userRequirementsFile

可选

指定用于初始化Head节点和Worker节点Python基础环境的requirements.txt文件。支持OSS和OSS-HDFS路径,路径格式必须为oss://<bucket>/<path>/requirements.txt。Ray集群启动后,系统自动执行pip install -r requirements.txt命令(异步执行,不影响集群启动)。