快速提交Ray任务

Updated at:

Ray is a general-purpose distributed computing framework widely used for distributed training, data processing, and model serving. Deep Learning Containers (DLC) of Platform for AI (PAI) allows you to submit Ray jobs directly. It automatically handles Ray cluster creation, job submission, and resource management, so you do not need to manually set up clusters or configure Kubernetes.

Prerequisites

If you use the SDK to submit a training job, you must configure environment variables. For more information, see Manage access credentials and Configure environment variables on Linux, macOS, and Windows.

Preparations

Prepare a node image

A Ray cluster includes at least two types of nodes: Head and Worker. DLC can use the same image for both node types. After you create a job, DLC automatically builds the Ray cluster and launches a Submitter node to submit the job to the cluster.

The Ray image version must be >= 2.6, and the image must include at least the components provided by ray[default]. The following images are supported:

  • PAI official images: PAI provides official images with Ray base components pre-installed.

    DLC platform images: In the image configuration dialog, set Framework to Ray. The following images are available:

    • ray:2.39.0-cpu-py312-ubuntu22.04 (3.90 GiB, CPU version)

    • ray:2.39.0-gpu-py312-cu118-ubuntu22.04 (13.61 GiB, GPU version)

  • Ray community images:

    • We recommend the Docker image rayproject/ray.

    • The rayproject/ray-ml image is also supported. This image includes built-in machine learning frameworks such as PyTorch and TensorFlow.

    When using GPUs, you must use a CUDA-enabled image. For more supported image versions, see the official Docker image documentation.

Prepare the startup command and script file

The startup command of a DLC job is used as the entrypoint command submitted by ray job submit. The startup command can be a single line or multiple lines. For example, python /root/code/sample.py, where:

  • sample.py is the Python script to run. You can mount the script file into the DLC container by using a dataset or code configuration. The following example shows the content of the script:

    import ray
    import os
    
    ray.init()
    
    @ray.remote
    class Counter:
        def __init__(self):
            # Used to verify runtimeEnv
            self.name = os.getenv("counter_name")
            # assert self.name == "ray"
            self.counter = 0
    
        def inc(self):
            self.counter += 1
    
        def get_counter(self):
            return "{} got {}".format(self.name, self.counter)
    
    counter = Counter.remote()
    
    for _ in range(50000):
        ray.get(counter.inc.remote())
        print(ray.get(counter.get_counter.remote()))
  • /root/code/ is the mount path.

Submit a training job

Submit a job in the console

  1. Go to the Create Job page.

    1. Log on to the PAI console. In the top navigation bar, select a region. On the right side, select a workspace, and then click Enter DLC.

    2. On the Distributed Training (DLC) page, click Create Job.

  2. On the Create Job page, configure the following key parameters. For more information about other parameters, see Create a training job.

    Parameter

    Description

    Example value

    Environment Information

    Node Image

    On the Alibaba Cloud Image tab, select a preset official Ray image.

    ray:2.39.0-cpu-py312-ubuntu22.04

    Startup Command

    The command submitted to the Ray cluster.

    python /root/code/sample.py

    Environment Variable

    Configure environment variables for Ray nodes through the runtime environment (runtime_env).

    Third-party Libraries

    You can configure a list of third-party libraries to set up the runtime environment dependencies (runtime_env) for Ray.

    Note

    In production environments, we strongly recommend that you use images with pre-installed dependencies to avoid job failures caused by installing dependency libraries at runtime.

    Not required

    Code Builds

    Upload your script file to the DLC container by using Online configuration or Local Upload.

    Use Local Upload:

    • Sample code file: sample.py.

    • Mount path: /root/code/.

    Resource Information

    Source

    Select Public Resources or Resource Quota.

    Public Resources

    Framework

    The framework type.

    Ray

    Job Resource

    • Node role:

      The Head node role is required, and at least one Worker role must be configured. You can add multiple Worker groups as needed.

    • Number of job nodes:

      The number of Head nodes is fixed at 1. By default, a Head node only runs the entrypoint script. Typically, you also need at least one Worker node to perform compute tasks. You can configure auto scaling to specify a range for the number of nodes. Each Ray job automatically generates a Submitter node to run the startup command. You can view the job status in the Submitter node logs.

    • Auto scaling:

      When you use Resource Quota, click the image icon to enable auto scaling for Worker nodes. After you set the maximum and minimum number of instances, the job enters the running state once the minimum required resources are met. The system then automatically adjusts the number of running instances based on resource availability and workload.

      Click More Configurations to specify the number of roles, GPU count, CPU count, and memory size. Then select an appropriate Scaling Policy. The following scaling policies are supported:

      • Default - maximize nodes: The system maximizes the number of nodes when sufficient resources are available.

      • RayAutoscaler - dynamic scaling: Supported only for Ray jobs. Ray Autoscaler automatically adjusts the number of nodes based on cluster load.

        Important

        To use the RayAutoscaler policy, you must bind a RAM role with read and write permissions on DLC jobs to the instance. In most cases, you can select the default PAI role.

      • CloudMonitorMetric - CloudMonitor metrics: The system automatically adjusts the number of nodes based on CloudMonitor metrics. You can configure target values for CPU utilization, memory utilization, GPU compute utilization, and GPU memory utilization.

    • Resource count:

      The logical resources on Ray cluster Worker nodes match the physical resources you configure when submitting the job. For example, if you configure one 8-GPU node, the Worker node in the Ray cluster has 8 GPUs of logical resources by default.

      To ensure cluster stability, the CPU logical resources on the Head node are set to 0 by default, and the Head node does not participate in compute scheduling. We recommend that the Head node has at least 2 GiB of memory, and you should increase this as the number of Tasks/Actors grows to avoid OOM errors. Resource configurations must match job requirements. We recommend using fewer large nodes instead of many small nodes.

      Note

      You can add the advanced parameter {"RayHeadScheduling": "true"} to enable the Head node to participate in compute scheduling.

    • Node role: Head and Worker

    • Number of nodes: 1 for each.

    • Instance type: Select ecs.g6.xlarge.

    Fault Tolerance and Diagnosis

    Head node fault tolerance

    You can configure a Redis instance for Head node fault tolerance. After you create a Redis instance, you must configure a VPC whitelist so that the DLC job VPC can access the Redis instance, and then configure a username and password.

    Important

    The default Redis username is default. Configuring a custom username is supported only in Ray 2.41 and later.

    Submitter retry count

    The maximum number of retries for the Submitter to submit a job to the Ray cluster.

    Ray log collection

    Before a Ray cluster is destroyed, you can persist the Ray cluster logs, framework logs, and task logs to OSS for later debugging and analysis.

    Important

    Ray log collection requires that you bind a role with read and write permissions on the target OSS bucket to the instance. Note: The default PAI role only has access to the default storage within the current workspace.

    You can download the log files directly from the OSS bucket later, or create a Ray History Server under AI Computing Asset Management>Jobs and configure the same storage path to view historical Ray cluster node logs and job status.

  3. After you configure the parameters, click OK.

Submit a job by using the SDK

  1. Install the Python DLC SDK.

    pip install alibabacloud_pai_dlc20201203==1.4.0
  2. Submit a DLC Ray job. The following example shows the sample code.

    #!/usr/bin/env python3
    
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_credentials.client import Client as CredClient
    
    from alibabacloud_pai_dlc20201203.client import Client as DLCClient
    from alibabacloud_pai_dlc20201203.models import CreateJobRequest
    
    region_id = '<region-id>'
    cred = CredClient()
    workspace_id = '12****'
    
    dlc_client = DLCClient(
        Config(credential=cred,
               region_id=region_id,
               endpoint='pai-dlc.{}.aliyuncs.com'.format(region_id),
               protocol='http'))
    
    create_job_resp = dlc_client.create_job(CreateJobRequest().from_map({
        'WorkspaceId': workspace_id,
        'DisplayName': 'dlc-ray-job',
        'JobType': 'RayJob',
        'JobSpecs': [
            {
                "Type": "Head",
                "Image": "dsw-registry-vpc.<region-id>.cr.aliyuncs.com/pai/ray:2.39.0-gpu-py312-cu118-ubuntu22.04",
                "PodCount": 1,
                "EcsSpec": 'ecs.c6.large',
            },
            {
                "Type": "Worker",
                "Image": "dsw-registry-vpc.<region-id>.cr.aliyuncs.com/pai/ray:2.39.0-gpu-py312-cu118-ubuntu22.04",
                "PodCount": 1,
                "EcsSpec": 'ecs.c6.large',
            },
        ],
        "UserCommand": "echo 'Prepare your ray job entrypoint here' && sleep 1800 && echo 'DONE'",
    }))
    job_id = create_job_resp.body.job_id
    print(f'jobId is {job_id}')

    Where:

    • region_id: The Alibaba Cloud region ID. For example, the region ID of China (Hangzhou) is cn-hangzhou.

    • workspace_id: The workspace ID. You can find it on the workspace details page. For more information, see Manage workspaces.

    • Image: Replace <region-id> with the actual Alibaba Cloud region ID. For example, the region ID of China (Hangzhou) is cn-hangzhou.

For more information about SDK usage, see Create a training job.

Submit a job by using the CLI

  1. Download the DLC client tool and complete user authentication. For more information, see Preparations.

  2. Submit a DLC Ray job. The following example shows the sample code.

    ./dlc submit rayjob --name=my_ray_job \
      --workers=1 \
      --worker_spec=ecs.g6.xlarge \
      --worker_image=dsw-registry-vpc.<region-id>.cr.aliyuncs.com/pai/ray:2.39.0-cpu-py312-ubuntu22.04 \
      --heads=1 \
      --head_image=dsw-registry-vpc.<region-id>.cr.aliyuncs.com/pai/ray:2.39.0-cpu-py312-ubuntu22.04 \
      --head_spec=ecs.g6.xlarge \
      --command="echo 'Prepare your ray job entrypoint here' && sleep 1800 && echo 'DONE'" \
      --workspace_id=4****

    For more information about configuring jobs submitted through the CLI, see Submit commands.

FAQ

Q: Why does my Ray job time out due to long environment preparation?

  • Check the logs of the Head node to verify whether the Ray runtime started successfully. If it did not, Ray is not available in the instance. Refer to the Preparations section to prepare a Ray-compatible image.

    After the Ray runtime starts successfully, the following output appears in the logs:

    -- --------------------
    -- Ray runtime started.
    -- --------------------
    
    -- Next steps
    -- To add another node to this Ray cluster, run
    --     ray start --address='192.168.128.203:6379'
    -- To connect to this Ray cluster:
    -- import ray
    -- ray.init()
    -- To terminate the Ray runtime, run
    --     ray stop
    -- To view the status of the cluster, use
    --     ray status
    -- --block
    -- This command will now block forever until terminated by a signal.
  • Check the event logs of the Head node. If the error "Readiness probe failed..." appears, it may indicate that the image is missing the dependencies required for the readiness check, or some transitive dependencies are unavailable. We recommend that you reinstall the ray[default] package in the original image by using pip or conda, or rebuild the image based on the official Ray image.

Appendix: Ray framework advanced configurations

Parameter

Type

Description

Default value

RayRuntimeEnv

string

Runtime environment dependencies

RayRedisAddress

string

External GCS Redis address

RayRedisUsername

string

External GCS Redis username. Supported only in Ray 2.41 and later

RayRedisPassword

string

External GCS Redis password

RaySubmitterBackoffLimit

int

Submitter retry count

0

RayObjectStoreMemoryBytes

int

Object store memory

30% of available memory

RayHeadScheduling

bool

Whether the Head node participates in scheduling

false

RayLogStoragePath

string

OSS bucket path for persisting Ray logs