Todos os produtos
Search
Central de documentação

E-MapReduce:Submit jobs to a Ray cluster

Última atualização: Jun 27, 2026

Os clusters Ray em workspaces do EMR Serverless Spark oferecem uma estrutura de computação distribuída nativa para Python, voltada ao treinamento e à inferência de modelos de machine learning. Crie e inicie um cluster Ray e envie jobs por meio de desenvolvimento interativo, SDK ou linha de comando.

Crie um cluster Ray

  1. Faça login no console do EMR Serverless Spark e acesse o workspace desejado.

  2. No painel de navegação à esquerda, clique em Cluster Management e, em seguida, na aba Ray clusters.

  3. Clique em Create Ray cluster. No painel exibido, configure os parâmetros abaixo e clique em Create.

    • Cluster name: insira um nome para o cluster.

    • Engine version: selecione uma versão do mecanismo. O padrão atual é emr-1.0.1 (Ray 2.47.1, Python 3.12).

    • Node groups: caso precise combinar recursos de CPU e GPU, crie grupos de nós separados para cada tipo. Clique em + Add worker node group e configure os seguintes parâmetros.

      • Node group name: defina um nome exclusivo para o grupo de nós.

      • Resource queue: escolha uma fila de recursos.

      • Resource specification: especifique a configuração de recursos dos nós.

      • Number of nodes: indique a quantidade de nós workers.

      • (Opcional) Network connection: se for necessário acessar serviços na mesma VPC a partir do cluster Ray, selecione uma conexão de rede existente.

      • (Opcional) Mount managed file directory: escolha um diretório de arquivos gerenciados existente para permitir que seus jobs Ray leiam e gravem arquivos nesse local. A montagem é um método comum de acesso a dados que mapeia caminhos do CPFS, NAS e OSS para um diretório local, possibilitando que seus programas acessem os dados via caminhos locais. Para criar um novo diretório de arquivos gerenciados, consulte Gerenciar diretórios de arquivos gerenciados.

      • (Opcional) Advanced cluster configurations: defina parâmetros avançados no formato JSON. Para mais detalhes, consulte Configurações avançadas do cluster.

    • Nota

      Atualmente, o console não oferece suporte à configuração de Auto Scaling para nós workers. Para configurar esse recurso, utilize a API. Para mais informações, consulte CreateRayCluster.

      Inicie o cluster Ray

      1. Na lista de clusters Ray, localize o cluster desejado e clique em Start.

      2. Aguarde até que o status do cluster mude para running.

        Nota

        A primeira inicialização de um cluster Ray em um workspace cria componentes adicionais e pode levar de 2 a 3 minutos. As inicializações subsequentes são significativamente mais rápidas.

      3. Após a inicialização do cluster, clique em invocation information para visualizar os endereços de endpoint e o token necessários ao envio de jobs Ray.

        Aviso

        Na versão atual, o token de um cluster Ray não expira. Mantenha seu token seguro para evitar acessos não autorizados.

        Clique em Dashboard para abrir a interface de monitoramento do cluster Ray e verificar o uso de recursos.

      Envie um job Ray

      Método 1: Desenvolvimento interativo

      Ideal para inícios rápidos e depuração. Conecte-se diretamente a um cluster Ray pelo ambiente Notebook no console do EMR Serverless Spark, sem configurações de rede adicionais. Para conectar-se a partir da máquina local, garanta conectividade de rede com o cluster Ray e um ambiente Python 3.12.

      1. Instale o cliente Ray:

        pip install ray[client]==2.47.1
      2. Na página invocation information do cluster, obtenha o gRPC address e o token.

      3. Use o código de exemplo abaixo para se conectar ao cluster Ray e realizar desenvolvimento interativo:

        import ray
        import os
        
        def get_metadata():
            headers = {"ray-token": "<your_token>"}
            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("Square results:", results)

        image

      Método 2: Envio via SDK

      Indicado para incorporar o envio de jobs Ray no código da sua aplicação. Seu ambiente local deve ter o Python 3.12 instalado.

      1. Instale a biblioteca cliente de envio de jobs Ray:

        pip install "ray[default]==2.47.1"
      2. Na página invocation information do cluster, obtenha o public endpoint e o token.

      3. Use o código de exemplo a seguir para se conectar ao cluster Ray e enviar um job:

        import time
        from ray.job_submission import JobSubmissionClient
        
        custom_headers = {
            "ray-token": "<your_token>"
        }
        
        client = JobSubmissionClient("<your_public_endpoint>", headers=custom_headers)
        
        # Alternatively, you can use the internal network. This requires the client to be in the same region and VPC as the cluster. 
        # Submitting over the internal network is recommended for improved stability.
        # 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. (Opcional) Para acessar arquivos em um diretório montado a partir do job, monte um diretório de arquivos gerenciados durante a criação do cluster Ray. Em seguida, referencie os arquivos pelo caminho de montagem ao enviar o job. O exemplo abaixo pressupõe a montagem de um diretório OSS no caminho /mnt/myoss:

        import time
        from ray.job_submission import JobSubmissionClient
        
        custom_headers = {
            "ray-token": "<your_token>"
        }
        
        client = JobSubmissionClient("<your_public_endpoint>", headers=custom_headers)
        
        # Use the mount path to directly reference a script file in OSS as the 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)
        Nota

        No código do job Ray, também é possível ler e gravar arquivos diretamente pelos caminhos de montagem. Por exemplo, leia de /mnt/myoss/test.txt ou grave um arquivo no diretório /mnt/myoss/. Os arquivos gravados nesse diretório são sincronizados com o caminho correspondente no OSS.

        O exemplo a seguir demonstra como ler e gravar arquivos em um diretório montado dentro de um job Ray:

        # The following code is submitted to the Ray cluster.
        
        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 'Write succeeded'
        
        print("--- Starting remote tasks ---")
        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)

      Método 3: Envio por linha de comando

      Ideal para integração com sistemas externos. Seu ambiente local deve ter o Python 3.12 instalado.

      1. Instale a biblioteca cliente de envio de jobs Ray:

        pip install "ray[default]==2.47.1"
      2. Na página invocation information do cluster, obtenha o public endpoint e o token.

      3. Execute os comandos abaixo para enviar um job Ray:

        # You must set the RAY_ADDRESS and RAY_JOB_HEADERS environment variables before you run `ray job submit`.
        export RAY_ADDRESS='<your_public_endpoint>'               
        export RAY_JOB_HEADERS='{"ray-token": "<your_token>"}'
        
        ## Submit the job
        ### Example: Submit a job that performs a distributed sort.
        ray job submit --working-dir "." -- python  test-ray-core-sort.py
        
        ### Supports reading from and writing to OSS.
        ray job submit --working-dir "." -- python test-saving-data-test.py
        
        ### Supports reading from and writing to OSS-HDFS.
        ray job submit --working-dir "." -- python test-saving-data-oss-hdfs-test.py
        
        ### The code supports reading from and writing to mounted directories.
        ray job submit --working-dir "." -- python  test-saving-data-test-mount.py

      Para mais parâmetros de linha de comando, consulte a documentação oficial da CLI de envio de jobs Ray.

      Configurações avançadas do cluster

      Especifique as configurações avançadas do cluster no formato JSON. Todos os parâmetros são opcionais.

      Parâmetro

      Obrigatório

      Descrição

      userDefinedFiles

      Opcional

      Especifica os arquivos do OSS a serem baixados nos nós head e worker durante a inicialização do cluster. Há suporte para caminhos OSS e OSS-HDFS. Separe múltiplos caminhos com vírgula (,). Os arquivos são baixados para o diretório /home/ray/work-dir em cada nó. Exemplo: oss://mybucket/hello.py,oss://mybucket2/test/test.jar

      userRequirementsFile

      Opcional

      Arquivo requirements.txt usado para inicializar o ambiente base Python nos nós head e worker. Há suporte para caminhos OSS e OSS-HDFS. O caminho deve estar no formato oss://<bucket>/<path>/requirements.txt. Após o início do cluster Ray, o sistema executa automaticamente pip install -r requirements.txt. Trata-se de uma operação assíncrona que não afeta a inicialização do cluster.