Ray on SURF ResearchCloud

Distrbuted computing on public infrastructure

Introduction

Ray.io is an open source framework to easily build distributed computing applications in Python.

ResearchCloud is a cloud platform offering compute resources on SURF infra, commercial cloud, or even your own OpenStack environment.

Why Ray?

Ray provides:

  • cluster: multiple head/worker nodes.
  • auto-scaler: creating/destroying nodes as work changes, so you don’t consume more resources than necessary. Mix CPU/GPU.
  • simple abstractions for distributed Python: turn functions into parallel tasks and classes into long-lived actors.
    • Plus: libraries for training, tuning, and serving.

Best suited to workloads that are:

  • Dynamic: new work is created or changed while the job runs.
  • Unpredictable: you can’t know in advance when or how much compute you’ll need.

Why Ray?

Examples:

  • Simulations, where each result decides which scenarios to run next.
  • Data pipelines: new data arriving decides what processing runs, and how much.
  • Model serving, when demand for inference fluctuates.
  • ML training and fine-tuning: Ray handles scaling, failures, and data feeding around frameworks like PyTorch.
  • Combinations: e.g. a pipeline that analyzes data, retrains a model, and serves it, all on one cluster.

Ray

import ray

ray.init()

@ray.remote(num_cpus=0.25)
def square(x):
    return x * x

futures = [square.remote(i) for i in range(2)]
results = ray.get(futures)
print(results)

Ray

import random
import ray

ray.init() # Connects to your existing cluster

@ray.remote(num_cpus=0.5)
class Player:
    def __init__(self, name: str):
        self.name = name

    def throw_hand(self) -> tuple:
        # Returns (Player Name, Choice)
        return self.name, random.choice(["Rock", "Paper", "Scissors"])

Ray

import gymnasium as gym
from ray.rllib.env.multi_agent_env import MultiAgentEnv

class RockPaperScissors(MultiAgentEnv):
    def __init__(self, config=None):
      ...
    def reset(self, *, seed=None, options=None):
      ...
    def step(self, action_dict):
     ...

config = (
    PPOConfig()
    .environment(RockPaperScissors, env_config={"max_steps": 10})
    # 1. Hardware Resource Specs Allocation
    .resources(
        num_cpus_for_main_process=1,       # Dedicates 1 CPU core to the coordinator/driver process
        num_gpus_for_main_process=0,       # Keeps the driver on the CPU
    )
    .env_runners(
        num_env_runners=2,
        num_cpus_per_env_runner=1,
        num_gpus_per_env_runner=0,
    )
    .learners(
        num_learners=1,
        num_cpus_per_learner=1,
        num_gpus_per_learner=0,
    )
    # 2. Multi-agent configuration mapping
    .multi_agent(
        policies={"p1": PolicySpec(), "p2": PolicySpec()},
        policy_mapping_fn=lambda agent_id, *a, **kw: "p1" if agent_id == "agent_1" else "p2",
    )
)

Ray on SRC

Ray already works out of the box on AWS, k8s, etc.

UU ITS are developing a Ray auto-scaler and catalog items for ResearchCloud.

Benefits:

  • Use SURF credits (NWO compute applications, RCCS).
  • Access to SRC GPUs (RTX6000 PRO).
  • Sovereign infrastructure on SURF Cloud.
    • or cheap and compliant use of commerical cloud
    • or in-house OpenStack

Ray on SRC

Steps:

  1. pip install "ray[default]" git+https://github.com/UtrechtUniversity/src-ray-provider.git
  2. export RESEARCH_CLOUD_TOKEN="..."
  3. edit cluster.yaml to specify your ResearchCloud CO/wallet.
  4. ray up cluster.yaml: cluster will be created on SRC.
  5. ray dashboard cluster.yaml: ssh-tunnel cluster to your local machine.

Now you are ready to run an application on the cluster, e.g.:

ray job submit --address http://localhost:8265 -- python3 test.py

Ray

cluster.yaml:

cluster_name: src-ray-example

provider:
  type: external
  module: src_ray_provider.node_provider.ResearchCloudNodeProvider
  co_name: ResearchCloud Development
  wallet_name: SRC account for ResearchCloud Development

Read more

  • Ray: A Distributed Framework for Emerging AI Applications (OSDI 2018): the original design paper https://arxiv.org/abs/1712.05889
  • RLlib: Abstractions for Distributed Reinforcement Learning https://arxiv.org/abs/1712.09381
  • Tune: A Research Platform for Distributed Model Selection and Training https://arxiv.org/abs/1807.05118
  • Exoshuffle: Large-Scale Shuffling at the Application Level https://arxiv.org/abs/2203.05072
  • Chrono::Ray: A Distributed Framework for High-Throughput Simulation-Based Analysis of Multibody Systems: Project Chrono + Ray for large simulation studies Paper: https://arxiv.org/abs/2605.13767 Code: https://github.com/uwsbel/chrono-ray