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
Faça login no console do EMR Serverless Spark e acesse o workspace desejado.
No painel de navegação à esquerda, clique em Cluster Management e, em seguida, na aba Ray clusters.
-
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.
Na lista de clusters Ray, localize o cluster desejado e clique em Start.
-
Aguarde até que o status do cluster mude para running.
NotaA 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.
-
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.
AvisoNa 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.
-
Instale o cliente Ray:
pip install ray[client]==2.47.1 Na página invocation information do cluster, obtenha o gRPC address e o token.
-
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)
-
Instale a biblioteca cliente de envio de jobs Ray:
pip install "ray[default]==2.47.1" Na página invocation information do cluster, obtenha o public endpoint e o token.
-
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)
-
(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)NotaNo 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.txtou 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) -
Instale a biblioteca cliente de envio de jobs Ray:
pip install "ray[default]==2.47.1" Na página invocation information do cluster, obtenha o public endpoint e o token.
-
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
NotaAtualmente, 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
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.
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.
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.
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
userDefinedFilesOpcional
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-direm cada nó. Exemplo:oss://mybucket/hello.py,oss://mybucket2/test/test.jaruserRequirementsFileOpcional
Arquivo
requirements.txtusado 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 formatooss://<bucket>/<path>/requirements.txt. Após o início do cluster Ray, o sistema executa automaticamentepip install -r requirements.txt. Trata-se de uma operação assíncrona que não afeta a inicialização do cluster.