Run A RayJob
This page shows how to leverage Kueue’s scheduling and resource management capabilities when running KubeRay’s RayJob.
This guide is for batch users that have a basic understanding of Kueue. For more information, see Kueue’s overview.
Before you begin
Make sure you are using Kueue v0.6.0 version or newer and KubeRay v1.1.0 or newer.
Check Administer cluster quotas for details on the initial Kueue setup.
See KubeRay Installation for installation and configuration details of KubeRay.
Note
In order to use RayJob, prior to v0.8.1, you need to restart Kueue after the installation. You can do it by running:kubectl delete pods -l control-plane=controller-manager -n kueue-system.RayJob definition
When running RayJobs on Kueue, take into consideration the following aspects:
a. Queue selection
The target local queue should be specified in the metadata.labels section of the RayJob configuration.
metadata:
labels:
kueue.x-k8s.io/queue-name: user-queue
b. Configure the resource needs
The resource needs of the workload can be configured in the spec.rayClusterSpec.
spec:
rayClusterSpec:
headGroupSpec:
template:
spec:
containers:
- resources:
requests:
cpu: "1"
workerGroupSpecs:
- template:
spec:
containers:
- resources:
requests:
cpu: "1"
c. Suspend control
Kueue controls the spec.suspend field of the RayJob. When a RayJob is admitted by Kueue, Kueue will unsuspend it by setting spec.suspend to false, regardless of its previous value.
d. Limitations
- A Kueue managed RayJob cannot use an existing RayCluster.
- The RayCluster should be deleted at the end of the job execution,
spec.ShutdownAfterJobFinishesshould betrue. - Because a Kueue workload can have a maximum of 18 PodSets, the maximum number of
spec.rayClusterSpec.workerGroupSpecsis 17.
Example RayJob
In this example, the code is provided to the Ray framework via a ConfigMap.
apiVersion: v1
kind: ConfigMap
metadata:
name: ray-job-code-sample
data:
sample_code.py: |
import ray
import os
import requests
ray.init()
@ray.remote
class Counter:
def __init__(self):
# Used to verify runtimeEnv
self.name = os.getenv("counter_name")
assert self.name == "test_counter"
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(5):
ray.get(counter.inc.remote())
print(ray.get(counter.get_counter.remote()))
# Verify that the correct runtime env was used for the job.
assert requests.__version__ == "2.26.0"The RayJob looks like the following:
apiVersion: ray.io/v1
kind: RayJob
metadata:
name: rayjob-sample
labels:
kueue.x-k8s.io/queue-name: user-queue
spec:
shutdownAfterJobFinishes: true
entrypoint: python /home/ray/samples/sample_code.py
runtimeEnvYAML: |
pip:
- requests==2.26.0
- pendulum==2.1.2
env_vars:
counter_name: "test_counter"
rayClusterSpec:
rayVersion: '2.55.1'
headGroupSpec:
rayStartParams:
dashboard-host: '0.0.0.0'
template:
spec:
containers:
- name: ray-head
image: rayproject/ray:2.55.1
ports:
- containerPort: 6379
name: gcs-server
- containerPort: 8265
name: dashboard
- containerPort: 10001
name: client
resources:
limits:
cpu: "1"
memory: "5Gi"
requests:
cpu: "1"
memory: "2Gi"
volumeMounts:
- mountPath: /home/ray/samples
name: code-sample
volumes:
- name: code-sample
configMap:
name: ray-job-code-sample
items:
- key: sample_code.py
path: sample_code.py
workerGroupSpecs:
- replicas: 1
minReplicas: 1
maxReplicas: 5
groupName: small-group
rayStartParams: {}
template:
spec:
containers:
- name: ray-worker
image: rayproject/ray:2.55.1
lifecycle:
preStop:
exec:
command: [ "/bin/sh","-c","ray stop" ]
resources:
limits:
cpu: "1"
requests:
cpu: "200m"
You can run this RayJob with the following commands:
# Create the code ConfigMap (once)
kubectl apply -f ray-job-code-sample.yaml
# Create a RayJob. You can run this command multiple times
# to observe the queueing and admission of the jobs.
kubectl create -f ray-job-sample.yaml
Autoscaling (a.k.a InTreeAutoscaling)
RayJob Autoscaling can automatically add or remove Ray worker pods based on resource demand.
This feature is supported since v0.15.2 and v0.14.7.
How to enable autoscaling in RayJob
- Enable feature gate for Elastic Workloads (Workload Slices)
ElasticJobsViaWorkloadSlices: true
- Add workload slicing annotation on RayJob
annotations:
kueue.x-k8s.io/elastic-job: "true"
- Enable
enableInTreeAutoscalingon RayJob
spec:
rayClusterSpec:
enableInTreeAutoscaling: true
Example RayJob with Autoscaling
In this example, the code is provided to the Ray framework via a ConfigMap.
apiVersion: v1
kind: ConfigMap
metadata:
name: ray-job-autoscaling-code-sample
data:
sample_code.py: |
import ray
import os
ray.init()
@ray.remote
def my_task(x, s):
import time
time.sleep(s)
return x * x
# run tasks in sequence to avoid triggering autoscaling in the beginning
print([ray.get(my_task.remote(i, 1)) for i in range(10)])
# run tasks in parallel to trigger autoscaling (scaling up)
print(ray.get([my_task.remote(i, 10) for i in range(10)]))
# run tasks in sequence to trigger scaling down
print([ray.get(my_task.remote(i, 1)) for i in range(30)])
The RayJob looks like the following:
apiVersion: ray.io/v1
kind: RayJob
metadata:
name: rayjob-autoscaling-sample
labels:
kueue.x-k8s.io/queue-name: user-queue
annotations:
kueue.x-k8s.io/elastic-job: "true"
spec:
shutdownAfterJobFinishes: true
entrypoint: python /home/ray/samples/sample_code.py
runtimeEnvYAML: |
pip:
- requests==2.26.0
- pendulum==2.1.2
env_vars:
counter_name: "test_counter"
rayClusterSpec:
rayVersion: '2.55.1'
autoscalerOptions:
idleTimeoutSeconds: 30
upscalingMode: Aggressive
enableInTreeAutoscaling: true
headGroupSpec:
rayStartParams:
dashboard-host: '0.0.0.0'
template:
spec:
containers:
- name: ray-head
image: rayproject/ray:2.55.1
ports:
- containerPort: 6379
name: gcs-server
- containerPort: 8265
name: dashboard
- containerPort: 10001
name: client
resources:
limits:
cpu: "1"
memory: "5Gi"
requests:
cpu: "1"
memory: "2Gi"
volumeMounts:
- mountPath: /home/ray/samples
name: code-sample
volumes:
- name: code-sample
configMap:
name: ray-job-autoscaling-code-sample
items:
- key: sample_code.py
path: sample_code.py
workerGroupSpecs:
- replicas: 1
minReplicas: 1
maxReplicas: 5
groupName: small-group
rayStartParams: {}
template:
spec:
containers:
- name: ray-worker
image: rayproject/ray:2.55.1
lifecycle:
preStop:
exec:
command: [ "/bin/sh","-c","ray stop" ]
resources:
limits:
cpu: "1"
requests:
cpu: "200m"
You can run this RayJob with the following commands:
# Create the code ConfigMap (once)
kubectl apply -f ray-job-autoscaling-code-sample.yaml
# Create a RayJob. You can run this command multiple times
# to observe the queueing and admission of the jobs.
kubectl create -f ray-job-autoscaling-sample.yaml
Feedback
Was this page helpful?
Glad to hear it! Please tell us how we can improve.
Sorry to hear that. Please tell us how we can improve.