ray-project/ray
Ray is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.
10 hidden assumptions · 10-stage pipeline · 8 components
Like any codebase, this repository makes assumptions it never checks — most are routine. The ones worth your attention are below, in plain language with what to do about each.
Runs distributed ML workloads across clusters by scheduling tasks, actors, and data pipelines
In the most common Ray Data → Ray Train pipeline: raw data (images, text, records) sits in cloud storage (S3, GCS). Ray Data reads it lazily in parallel — each file becomes one or more Arrow RecordBatch blocks, each block an ObjectRef in the Plasma store. The dataset passes through a chain of map_batches() operators: preprocessing (resize, normalize), then GPU inference or feature extraction. The streaming executor drives these operators concurrently with backpressure so the cluster's object store memory doesn't overflow. The transformed blocks are either written back to cloud storage or iterated by Ray Train workers as training batches. Ray Train workers each see a shard of the data, run forward+backward passes, synchronize gradients via NCCL all-reduce, and checkpoint the model to shared storage. In the serving path, HTTP requests arrive at a Ray Serve HTTP proxy actor, are routed to a replica actor, the model runs inference, and the response is returned over HTTP.
Under the hood, the system uses 5 feedback loops, 5 data pools, 5 control points to manage its runtime behavior.
A 8-component repository. 7374 files analyzed. Data flows through 10 distinct pipeline stages.
Hidden Assumptions
Most of what this code assumes is routine. These 2 are the ones most likely to cause trouble here — in plain terms, with what to do about each. The rest are minor; they're under "Show everything".
These benchmark scripts assume they have permission to read from and write to specific internal Amazon S3 buckets before they start any work. There is no early check. If your credentials are wrong, expired, or missing, the job will spin up an entire cluster, process data for potentially hours, and then fail right at the end when it tries to save results — with a confusing error message.
What to do: Before kicking off a long run, do a quick manual check that you can list and write to the target storage buckets from the machine or role that will run the job.
The image embedding job has the batch size permanently set to 1024 images, tuned for one specific type of GPU (an A10G). If you run it on a different GPU — or on a regular CPU machine without a special fake-GPU label configured — it either crashes with an out-of-memory error immediately or hangs forever waiting for hardware that isn't there, with no helpful message explaining why.
What to do: Before running the embedding benchmark, confirm your cluster's GPU type matches what the script expects, or adjust the batch size to one that fits your hardware.
Show everything (8 more)
The benchmark assumes every image in the dataset is exactly the same large size. The block count and memory estimates are all calculated from this assumption. If the actual images are a different size, the job's memory usage and speed estimates will be wrong — it might run out of memory silently or report misleading throughput numbers.
What to do: If you're pointing this benchmark at your own dataset, verify that the typical image dimensions match what the script was written for, or update the block-size constant to match your data.
release/nightly_tests/dataset/image_embedding_from_uris/main.py:create_metadata
The text embedding job silently tries to fetch a private API key from a specific Amazon secrets service when each worker starts up. If the worker machines don't have permission to access that secret — which is easy to miss in a new environment — every worker quietly fails to start, and the whole job hangs with no clear explanation of why.
What to do: Make sure the machines running this job have been granted access to the specific secrets entry it needs, and test that access before starting a long run.
release/nightly_tests/dataset/text_embedding/main.py:TextEmbedder.__init__
The training ingest benchmark fills up the cluster's memory on purpose to test backpressure, but it never checks whether the cluster actually has enough memory configured for this test to be meaningful. On a smaller cluster, the test silently spills data to disk and reports artificially slow timings that look like regressions but are really just a memory configuration issue.
What to do: Before treating a slow benchmark result as a real regression, confirm the cluster's memory configuration matches what the test was designed for.
release/nightly_tests/dataset/training_ingest_regression_test/main.py:train_loop_per_worker
The video analysis app accepts a storage path in every incoming request and immediately tries to download from it, without checking whether that path is actually in the expected storage bucket or whether the server has permission. A misconfigured server or a caller pointing at a bucket the server can't reach gets a confusing internal error instead of a clear explanation.
What to do: Add a check that the storage path in each request matches the bucket the server is configured to access, and return a clear error message if it doesn't.
doc/source/serve/tutorials/video-analysis/app.py:VideoAnalyzer.analyze
When you ask Ray to install a specific set of packages on worker machines, it caches the result so it doesn't reinstall every time. But if you use version ranges (like 'give me version 2 or newer') instead of exact versions, the cache won't notice when a newer version of that package is released. Workers will quietly keep using the old version, and your results may silently change or be wrong without any warning.
What to do: Pin all package versions to exact numbers in your runtime environment configuration rather than using ranges, so the cache behaves predictably.
python/ray/_private/runtime_env/agent/runtime_env_agent.py:RuntimeEnvAgent
The monitoring dashboard always shows information that is slightly out of date — it refreshes on a timer rather than updating instantly. During a fast-moving incident (a node crashing, a job finishing), what you see on screen may be several seconds behind reality, which can cause confusion if you're trying to react quickly.
What to do: When debugging a live issue, rely on the Ray command-line tools for up-to-the-moment state rather than trusting the dashboard display.
python/ray/dashboard/client/src/App.tsx:App
The image preprocessing settings (how pixel values are scaled and centered before going into the model) are hardcoded numbers that only work correctly for one specific model. If someone swaps in a different model variant, the preprocessing will be silently wrong — the model will still run and produce numbers, but those numbers will be meaningless, and any downstream use of the embeddings (like similarity search) will give wrong answers.
What to do: Load the preprocessing configuration from the model itself rather than hardcoding the normalization values, so they automatically match whatever model is in use.
release/nightly_tests/dataset/image_embedding_from_jsonl/main.py:ImageEmbedder.__call__
You can tell the benchmark to run many parallel GPU workers, but nothing stops you from requesting more workers than you have GPUs. If that happens, multiple workers end up sharing the same GPU, each trying to load a large model, and they crash each other out of memory — then keep retrying in a loop that burns cluster time without getting anything done.
What to do: Set the maximum number of parallel inference workers to be no greater than the number of GPUs actually available in your cluster.
release/nightly_tests/dataset/image_embedding_from_uris/main.py:parse_args
Open the standalone hidden-assumptions report for ray →
How Data Flows Through the System
In the most common Ray Data → Ray Train pipeline: raw data (images, text, records) sits in cloud storage (S3, GCS). Ray Data reads it lazily in parallel — each file becomes one or more Arrow RecordBatch blocks, each block an ObjectRef in the Plasma store. The dataset passes through a chain of map_batches() operators: preprocessing (resize, normalize), then GPU inference or feature extraction. The streaming executor drives these operators concurrently with backpressure so the cluster's object store memory doesn't overflow. The transformed blocks are either written back to cloud storage or iterated by Ray Train workers as training batches. Ray Train workers each see a shard of the data, run forward+backward passes, synchronize gradients via NCCL all-reduce, and checkpoint the model to shared storage. In the serving path, HTTP requests arrive at a Ray Serve HTTP proxy actor, are routed to a replica actor, the model runs inference, and the response is returned over HTTP.
- Read raw data from cloud storage into Arrow blocks — ray.data.read_parquet(), read_images(), or read_json() creates a Dataset with a ReadOperator as the root of the logical plan. The streaming executor launches Ray tasks (one per file or file chunk) that read the file bytes, parse them into pyarrow.Table objects, and store each table as an ObjectRef in the Plasma store. Block size is controlled by target_max_block_size. For the image benchmarks, each block holds ~11 images (IMAGES_PER_BLOCK = 11 in image_embedding_from_uris/main.py). (config: ray.defaults.python, ray.platforms)
- Preprocess data with map_batches() operators — User-defined transform functions (e.g. preprocess() in image_classification_from_parquet/main.py which applies torchvision ResNet50_Weights.DEFAULT.transforms()) are wrapped in a MapBatchesOperator. The streaming executor submits these as Ray tasks or actor-pool tasks (ActorPoolStrategy). Each invocation receives a DataBatch (Dict[str, np.ndarray]), applies the transform, and returns a new DataBatch. For the DreamBooth dataset (release/air_examples/dreambooth/dreambooth/dataset.py), this includes tokenization via AutoTokenizer and image transforms via torchvision.transforms.Compose. [DataBatch → DataBatch]
- Buffer and backpressure between pipeline stages — StandardStreamingExecutor (python/ray/data/_internal/execution/streaming_executor.py) maintains per-operator output queues of ObjectRefs. When a downstream operator's queue exceeds its memory budget, the executor pauses submitting new tasks to the upstream operator. This prevents the object store from filling up when a slow GPU inference step can't keep up with a fast preprocessing step. The budget is derived from the total object store size and the number of active operators. [DataBatch → DataBatch]
- GPU inference via actor pool (map_batches with GPU actors) — For workloads like image embedding (image_embedding_from_jsonl/main.py), a class with a __call__ method loads a model (ViTForImageClassification) in its __init__ and processes batches. Ray Data keeps a pool of these actors alive (concurrency controlled by --inference-concurrency min/max args), dispatching blocks to idle actors. Each actor receives a DataBatch of decoded images (shape [B, 224, 224, 3] after ViTImageProcessor) and returns embeddings. BATCH_SIZE = 1024 images per call. [DataBatch → DataBatch]
- Stream processed blocks into Ray Train workers — When a Ray Dataset is passed to TorchTrainer, each training worker calls session.get_dataset_shard('train') which returns an iterator over the worker's assigned data shard. Internally, Ray Data streams blocks from the object store to the worker, converting each pyarrow.Table block to a torch.Tensor batch via a collation function. In training_ingest_regression_test/main.py, BATCH_SIZE = 1024 images per step; time-to-first-batch and per-batch latency are measured to detect pipeline regressions. [DataBatch → DataBatch]
- Distributed training forward + backward pass — Each TorchTrainer worker runs the user's train_loop_per_worker function. It wraps the model with ray.train.torch.prepare_model() (which applies PyTorch DDP wrapping for multi-GPU synchronization), runs forward pass, computes loss, calls loss.backward(), and optimizer.step(). Gradients are all-reduced across workers via NCCL. In the ingest benchmark, the model is a tiny CNN (Conv2d → AdaptiveAvgPool2d → Linear) so compute is cheap and the data pipeline is the bottleneck. [DataBatch]
- Checkpoint model state to shared storage — At the end of each epoch (or configurable interval), TorchTrainer calls ray.train.report(metrics, checkpoint=Checkpoint.from_dict({'model': state_dict})). The checkpoint is serialized and written to the configured storage path (local shared filesystem or cloud storage). In Tune-wrapped training, the best checkpoint is selected by the Tuner based on the reported metric (e.g. lowest validation loss).
- Write output blocks to cloud storage — ray.data.Dataset.write_parquet() or write_json() at the end of an inference pipeline submits write tasks that serialize each Arrow block to Parquet format and upload to the destination (e.g. WRITE_PATH = 's3://ray-data-write-benchmark/{uuid}' in the benchmark scripts). Write is parallelized — one task per output block. [DataBatch]
- HTTP request enters Ray Serve and routes to replica — An incoming HTTP POST (e.g. to /analyze in the video-analysis app.py) hits the Ray Serve HTTP Proxy actor (an aiohttp server). The proxy deserializes the JSON body into an AnalyzeRequest Pydantic model (stream_id: str, video_path: str S3 URI, num_frames: int, chunk_duration: float, use_batching: bool). It looks up the matching deployment in the route table and sends the request to an available replica actor via a DeploymentHandle. The replica runs the user code and returns the response. [ServeDeploymentConfig]
- Autoscaler reconciles cluster size to resource demand — StandardAutoscaler.update() (python/ray/autoscaler/_private/autoscaler.py) runs every few seconds on the head node. It queries GCS for pending resource demands (tasks waiting for CPUs/GPUs that don't exist yet), computes how many nodes of each configured type (from the cluster YAML's available_node_types) would satisfy the demand, and calls NodeProvider.create_node() for the delta. It also terminates nodes that have been idle (no running tasks) for longer than idle_timeout_minutes. The result is a new ClusterAutoscalerState. [ClusterAutoscalerState → ClusterAutoscalerState] (config: ray.defaults.python, ray.architectures, ray.platforms)
Data Models
The data structures that flow between stages — the contracts that hold the system together.
python/ray/_raylet.pyxAn opaque 20-byte identifier (object ID + owner address) that points to a value stored in the distributed object store (Plasma). Passed between tasks as futures — calling ray.get(ref) blocks until the value is available and returns the deserialized Python object.
Created when a task is submitted (ray.remote call) or ray.put() is called; stored in Plasma until all references are released; garbage collected when ref-count drops to zero.
src/ray/core_worker/task_manager.hProtobuf message with fields: task_id (bytes), function_descriptor (language + module + class + function name), args (list of ObjectRef or serialized values), num_returns (int), required_resources (dict: CPU/GPU/memory -> float), scheduling_strategy (placement group or node affinity)
Constructed by CoreWorker on task submission, sent to the local Raylet via gRPC, stored in the task queue until a worker with matching resources is available, then dispatched to that worker.
python/ray/data/dataset.pyA lazy computation graph (DAG of LogicalOperator nodes) over Arrow RecordBatches. Each node declares its input schema, output schema, and a transform function. Not materialized until .materialize() or iteration begins. Internally, blocks are ObjectRefs pointing to pyarrow.Table objects in the object store.
Created by ray.data.read_*() calls; operators are chained lazily; execution begins when consumed (iteration, write, materialize); blocks flow through the streaming executor stage by stage.
release/nightly_tests/dataset/image_classification_from_parquet/main.pyDict[str, np.ndarray] where keys are column names (e.g. 'image': ndarray of shape [B, H, W, C], 'label': ndarray of shape [B]). This is the unit of data flowing through Ray Data map_batches() operators.
Deserialized from a pyarrow.Table block, passed into user-defined transform functions (preprocessing, inference), then serialized back to Arrow and stored in the object store as the next stage's input block.
python/ray/_private/runtime_env/agent/runtime_env_agent.pyDict with optional keys: pip (list of package strings), conda (str or dict), container (dict with image/run_options), env_vars (dict), working_dir (str URI), py_modules (list of URIs). Attached to a job, actor class, or remote function.
Specified by user at job/actor submission time; serialized and sent to the RuntimeEnvAgent on the target node via gRPC; agent installs packages, caches the result, then reports success/failure back before the worker process starts.
python/ray/_private/runtime_env/agent/runtime_env_agent.pyDataclass with fields: success (bool), result (str — path to the env or error message), creation_time_ms (int)
Created by RuntimeEnvAgent after attempting to set up the environment; returned to the Raylet to unblock the waiting worker launch.
python/ray/serve/deployment.pyPydantic model with fields: num_replicas (int), ray_actor_options (dict: resources per replica), max_ongoing_requests (int), autoscaling_config (AutoscalingConfig with min/max replicas, target_ongoing_requests), route_prefix (str), graceful_shutdown_timeout_s (float)
Defined by user via @serve.deployment decorator or config file; sent to ServeController actor which reconciles desired vs actual replica count by creating/destroying Ray actors.
python/ray/autoscaler/_private/autoscaler.pyDataclass with fields: active_nodes (Dict[NodeType, int]), idle_nodes (Dict[NodeType, int]), pending_nodes (List[Tuple[NodeIP, NodeType, NodeStatus]]), pending_launches (Dict[NodeType, int]), failed_nodes (List[Tuple[NodeIP, NodeType]]), node_availability_summary (NodeAvailabilitySummary), pending_resources (Dict[str, int])
Rebuilt every autoscaler tick by querying GCS for current node states and resource demands; used to decide how many nodes of each type to add or remove; published to the dashboard for display.
System Behavior
How the system operates at runtime — where data accumulates, what loops, what waits, and what controls what.
Data Pools
A shared-memory store on each node that holds serialized Ray objects (task arguments, return values, Dataset blocks). Workers on the same node read objects zero-copy via memory-mapped files. Objects are reference-counted; when all ObjectRefs are released, the object is eligible for eviction. When full, Ray spills objects to disk (a configurable spill directory).
Cluster-wide authoritative record of all actors (location, state), nodes (IP, resources, liveness), jobs (status, runtime env), and placement groups. Every component that needs to find where an actor lives or register a new one goes through GCS. Backed by an in-process store (or Redis for HA deployments).
A per-node directory managed by RuntimeEnvAgent that caches installed environments keyed by a hash of the RuntimeEnvConfig spec. On cache hit, the agent returns immediately without reinstalling. On cache miss, it installs (pip, conda, or docker) and stores the result. Prevents re-downloading packages when multiple actors on the same node share an environment.
Persistent storage (local path or cloud URI like s3://) where Ray Train writes model checkpoints. Each checkpoint is a directory containing model weights and optimizer state. Ray Tune reads these to resume trials after failures or to select the best trial result.
The dashboard backend polls GCS and Ray State APIs every API_REFRESH_INTERVAL_MS milliseconds and caches the results in memory. The React frontend polls these cached results. This means dashboard state is eventually consistent — newly created actors or completed tasks appear with up to one refresh interval of delay.
Feedback Loops
- Autoscaler scale-up/scale-down loop (polling, balancing) — Trigger: Timer fires every ~5 seconds on head node. Action: StandardAutoscaler.update() reads pending resource demands from GCS, computes desired node count per type, launches or terminates nodes via cloud provider API, updates ClusterAutoscalerState. Exit: Cluster reaches steady state where pending demands == 0 and no idle nodes exceed idle_timeout_minutes.
- Ray Serve replica autoscaling loop (auto-scale, balancing) — Trigger: ServeController's autoscaling goroutine polls ongoing request counts from each replica at a configurable interval. Action: If average ongoing_requests per replica exceeds target_ongoing_requests in AutoscalingConfig, create new replica actors; if below, send shutdown signals to excess replicas (down to min_replicas). Exit: Replica count stabilizes within [min_replicas, max_replicas] at the target load.
- Ray Data streaming backpressure loop (backpressure, balancing) — Trigger: Downstream operator's output queue (in ObjectRefs) exceeds its memory budget. Action: StandardStreamingExecutor stops submitting new tasks to the upstream operator until the downstream queue drains below threshold. Exit: Downstream queue drops below the low-watermark, upstream resumes.
- RLlib training loop (rollout → learn → update) (training-loop, reinforcing) — Trigger: Algorithm.train() called by user or Tune. Action: Rollout workers (Ray actors running environment episodes) collect experience trajectories; learner actor consumes trajectories to compute policy gradient and update weights; updated weights are broadcast back to rollout workers. Exit: Configured number of training iterations or convergence criterion met.
- Task retry on worker failure (retry, balancing) — Trigger: A worker process dies (SIGKILL, OOM, node failure) while executing a task. Action: Raylet/GCS detects the worker death, marks the task as failed, re-submits it (up to max_retries times, default 3) to another available worker. Exit: Task succeeds or max_retries exhausted (raises RayTaskError to caller).
Delays
- Runtime environment installation (async-processing, ~Seconds to minutes depending on package set size) — Worker processes are blocked (not accepting tasks) until RuntimeEnvAgent reports success. On cache hit this is near-zero; on miss for large pip environments it can take minutes, blocking task execution.
- Dashboard data refresh (eventual-consistency, ~API_REFRESH_INTERVAL_MS (constant in python/ray/dashboard/client/src/common/constants.ts)) — Dashboard UI shows cluster state that is up to one refresh interval stale. Newly launched actors, completed tasks, or node failures appear with this delay.
- Ray Data object store spilling (async-processing, ~Variable — depends on disk I/O speed) — When the Plasma store is full, objects are spilled to a local disk directory asynchronously. Downstream tasks that need a spilled object must wait for it to be read back from disk, introducing pipeline stalls.
- Benchmark step sleep (backpressure simulation) (batch-window, ~--step-sleep-s parameter (e.g. 2.0 seconds)) — In the peak_object_store_memory test configuration, each training step sleeps for step-sleep-s seconds to simulate a slow GPU computation. This backs up the Ray Data pipeline, filling object store queues and exercising the backpressure mechanism.
Control Points
- ray.defaults.python / ray.python versions (architecture-switch) — Controls: Which Python versions Docker images are built for (3.10, 3.11, 3.12, 3.13, 3.14) and which is the default (3.10). Determines the Python interpreter used in all Ray worker processes.. Default: default: 3.10
- ray.platforms (GPU/CPU platform) (device-selection) — Controls: Which CUDA versions Docker images are built for (cpu, tpu, cu11.7.1-cudnn8 through cu13.0.0-cudnn). Determines whether GPU workloads can run and which CUDA toolkit version is available.. Default: cpu, tpu, cu11.7.1-cudnn8 through cu13.0.0-cudnn
- --inference-concurrency (min, max) (hyperparameter) — Controls: The minimum and maximum number of concurrent GPU inference actor instances Ray Data maintains in its actor pool for the ViT inference operator. Higher values increase GPU utilization but require more GPU memory.. Default: Required CLI argument, no default
- use_batching (Serve request batching) (feature-flag) — Controls: Whether Ray Serve batches multiple incoming HTTP requests together before passing them to the model. When True, multiple /analyze requests arriving within a time window are grouped into one batch call, improving GPU utilization at the cost of added latency for early-arriving requests.. Default: False (default)
- RAY_LOGGING_CONFIG / LoggingConfig (runtime-toggle) — Controls: Per-worker log encoding (TEXT or JSON), log level (DEBUG/INFO/WARNING/ERROR), and which additional standard log attributes to include. Affects all Ray worker, actor, and driver processes.. Default: encoding: TEXT, log_level: INFO
Technology Stack
Implements the performance-critical core: Raylet task scheduling, Plasma object store shared memory management, GCS metadata service, and CoreWorker distributed primitives. All cross-process communication uses gRPC.
The primary user-facing language. All AI library APIs (Data, Train, Tune, Serve, RLlib) are Python. The _raylet.pyx Cython extension bridges Python to the C++ CoreWorker.
python/ray/_raylet.pyx is a Cython extension that wraps the C++ CoreWorker class, giving Python code zero-overhead access to task submission, object store operations, and actor management.
The in-memory columnar data format for all Ray Data blocks. Arrow's zero-copy IPC format means blocks can be passed between workers via shared memory (Plasma) without serialization overhead. All Ray Data transforms operate on pyarrow.Table objects.
The Ray Dashboard is a React SPA (python/ray/dashboard/client/src/). Axios (via requestHandlers.ts) polls the Python dashboard backend for cluster state; Material UI renders it. The app handles authentication token flows for secured clusters.
Used by RuntimeEnvAgent (per-node environment installer) as its HTTP server to receive environment creation requests from the Raylet. Also used by the dashboard backend for its REST API.
Defines the wire format for all inter-component messages (TaskSpec, GCS RPCs, runtime env requests). Generated stubs live in python/ray/core/generated/ and src/ray/protobuf/.
Used by Ray Train workers for model training (DDP, FSDP) and by benchmark inference actors (ResNet50, ViT). Ray Train's torch integration handles DDP setup, device placement, and gradient synchronization.
Key Components
- CoreWorker (C++) (orchestrator) — The central per-worker runtime object. Every Python worker process (and actor) has exactly one CoreWorker. It submits tasks to the local Raylet, tracks ObjectRef lifetimes via reference counting, manages actor creation/deletion, serializes arguments into the object store, and deserializes results. The Python _raylet.pyx Cython extension exposes CoreWorker's methods to Python as ray.remote(), ray.get(), ray.put().
src/ray/core_worker/core_worker.cc - Raylet (scheduler) — A C++ process running on every cluster node that accepts tasks from CoreWorker, manages a pool of Python worker processes, matches tasks to workers based on resource availability (CPU, GPU, memory, custom resources), and fetches any missing object arguments from remote nodes' object stores before dispatching. It also reports resource usage to GCS so the autoscaler can make scaling decisions.
src/ray/raylet/ - GcsClient / GCS Server (store) — Global Control Store — the cluster's authoritative metadata service. Stores actor locations, node liveness, job records, resource totals, and placement group state. All other components (Raylet, CoreWorker, dashboard, autoscaler) query GCS to find where actors live or to register new objects. Backed by Redis in older deployments, now often by an in-process store.
src/ray/gcs/ - StandardStreamingExecutor (executor) — The engine that actually runs a Ray Data pipeline. It takes the logical DAG of operators (map, filter, read, write) and drives them concurrently: for each operator, it launches Ray tasks or actor-pool tasks to process blocks, applies backpressure (pausing upstream operators when downstream queues are full) to bound memory usage, and streams completed blocks to the next operator without waiting for the full dataset to materialize.
python/ray/data/_internal/execution/streaming_executor.py - ServeController (orchestrator) — A long-running Ray actor (singleton per Serve cluster) that owns the desired state of all deployments. When a user calls serve.run() or updates a config, ServeController reconciles: it diffs the desired deployment config against the current replica count, creates new replica actors or sends shutdown signals to excess ones, and updates the routing table so the HTTP proxy knows where to send requests.
python/ray/serve/_private/controller.py - StandardAutoscaler (optimizer) — Runs on the head node and periodically (every few seconds) computes the desired cluster size. It reads pending resource demands from GCS (tasks waiting for resources), computes how many nodes of each configured type would satisfy those demands, and calls the cloud provider API (AWS, GCP, Azure, K8s) to launch or terminate nodes. It also enforces idle timeout — terminating nodes that have had no work for a configurable period.
python/ray/autoscaler/_private/autoscaler.py - RuntimeEnvAgent (executor) — An aiohttp HTTP server running as a subprocess on every node. When a worker needs a runtime environment (e.g. specific pip packages), the Raylet calls the local RuntimeEnvAgent via HTTP. The agent checks its cache; on a miss, it installs the packages (pip install, conda create, or docker pull), records the result in a local cache keyed by the env spec hash, and reports back. The creation result (RuntimeEnvCreationResult) unblocks the waiting worker.
python/ray/_private/runtime_env/agent/runtime_env_agent.py - TorchTrainer (orchestrator) — User-facing entry point for distributed PyTorch training. When .fit() is called, it launches N Ray actors (one per worker, sized by ScalingConfig specifying num_workers and resources_per_worker), sets up PyTorch distributed (initializes NCCL process group, assigns ranks), and runs the user's training function on each worker in parallel. It handles checkpoint saving/restoring and reports metrics back to Ray Tune if used within a hyperparameter search.
python/ray/train/torch/torch_trainer.py
Package Structure
The main Ray Python package containing the core distributed runtime, all AI libraries (Data, Train, Tune, Serve, RLlib), the dashboard web UI, and the autoscaler. This is what users install via pip.
Release testing infrastructure and benchmark suite that exercises Ray's AI libraries at scale, including nightly regression tests for Ray Data pipelines, training throughput, serving latency, and LLM workloads.
Explore the interactive analysis
See the full architecture map, data flow, and code patterns visualization.
Analyze on CodeSeaRelated Repository Repositories
Frequently Asked Questions
What is ray used for?
Runs distributed ML workloads across clusters by scheduling tasks, actors, and data pipelines ray-project/ray is a 8-component repository written in Python. Data flows through 10 distinct pipeline stages. The codebase contains 7374 files.
How is ray architected?
ray is organized into 9 architecture layers: C++ Core Runtime, Python Runtime Bridge, Ray Data (Distributed ETL/Inference Pipeline), Ray Train (Distributed Model Training), and 5 more. Data flows through 10 distinct pipeline stages. This layered structure keeps concerns separated and modules independent.
How does data flow through ray?
Data moves through 10 stages: Read raw data from cloud storage into Arrow blocks → Preprocess data with map_batches() operators → Buffer and backpressure between pipeline stages → GPU inference via actor pool (map_batches with GPU actors) → Stream processed blocks into Ray Train workers → .... In the most common Ray Data → Ray Train pipeline: raw data (images, text, records) sits in cloud storage (S3, GCS). Ray Data reads it lazily in parallel — each file becomes one or more Arrow RecordBatch blocks, each block an ObjectRef in the Plasma store. The dataset passes through a chain of map_batches() operators: preprocessing (resize, normalize), then GPU inference or feature extraction. The streaming executor drives these operators concurrently with backpressure so the cluster's object store memory doesn't overflow. The transformed blocks are either written back to cloud storage or iterated by Ray Train workers as training batches. Ray Train workers each see a shard of the data, run forward+backward passes, synchronize gradients via NCCL all-reduce, and checkpoint the model to shared storage. In the serving path, HTTP requests arrive at a Ray Serve HTTP proxy actor, are routed to a replica actor, the model runs inference, and the response is returned over HTTP. This pipeline design reflects a complex multi-stage processing system.
What technologies does ray use?
The core stack includes C++ (with Abseil, Boost, gRPC) (Implements the performance-critical core: Raylet task scheduling, Plasma object store shared memory management, GCS metadata service, and CoreWorker distributed primitives. All cross-process communication uses gRPC.), Python (CPython) (The primary user-facing language. All AI library APIs (Data, Train, Tune, Serve, RLlib) are Python. The _raylet.pyx Cython extension bridges Python to the C++ CoreWorker.), Cython (python/ray/_raylet.pyx is a Cython extension that wraps the C++ CoreWorker class, giving Python code zero-overhead access to task submission, object store operations, and actor management.), Apache Arrow / PyArrow (The in-memory columnar data format for all Ray Data blocks. Arrow's zero-copy IPC format means blocks can be passed between workers via shared memory (Plasma) without serialization overhead. All Ray Data transforms operate on pyarrow.Table objects.), React + TypeScript (with Material UI, Axios) (The Ray Dashboard is a React SPA (python/ray/dashboard/client/src/). Axios (via requestHandlers.ts) polls the Python dashboard backend for cluster state; Material UI renders it. The app handles authentication token flows for secured clusters.), aiohttp (Used by RuntimeEnvAgent (per-node environment installer) as its HTTP server to receive environment creation requests from the Raylet. Also used by the dashboard backend for its REST API.), and 2 more. A focused set of dependencies that keeps the build manageable.
What system dynamics does ray have?
ray exhibits 5 data pools (Plasma Object Store, GCS Metadata Store), 5 feedback loops, 5 control points, 4 delays. The feedback loops handle polling and auto-scale. These runtime behaviors shape how the system responds to load, failures, and configuration changes.
What design patterns does ray use?
5 design patterns detected: Actor-as-Service, Lazy Computation DAG with Streaming Execution, Protobuf-over-gRPC for All Internal Communication, Benchmark-as-a-Release-Gate, Pydantic for API Contract Validation at Serve Boundaries.
Analyzed on September 10, 2026 by CodeSea. Written by Karolina Sarna.