RAY GPU ARCHITECTURE: PLASMA OBJECT STORE AND GPU MEMORY MANAGEMENT
Ray's GPU architecture centers on the Plasma object store and GPU-to-object-store data transfer model. When a Ray task runs on a GPU, the tensors reside in GPU memory and are referenced by Ray's distributed object store via a plasma ID. The challenge is that Ray's default object store (shared memory at `/dev/shm`) copies tensors through CPU RAM before GPU access. For a 16 GB GPU tensor, the transfer path is: GPU HBM -> CPU DRAM (via `cudaMemcpy`) -> Plasma shared memory (via DRAM copy) -> CPU DRAM -> GPU HBM (for the consuming task). This double-copy overhead adds 50-150ms per 16 GB transfer on PCIe Gen5 (32 GB/s), or 8-25ms on NVLink (900 GB/s). Ray 2.30+ addresses this with `ray.init(gpu_direct=True)` which enables GPUDirect RDMA between GPUs and Ray's shared memory, bypassing CPU copies entirely.
The GPU memory manager in Ray 2.33+ (`RAY_gpu_memory_manager=1`) introduces explicit GPU memory budgeting per worker. Each Ray actor with `num_gpus=1` receives a memory quota of 70% of VRAM by default (56 GB on 80 GB H100), with the remaining 30% reserved for Ray's internal GPU objects and model loading overhead. The `@ray.remote(num_gpus=1, max_gpu_memory=60*1024**3)` decorator sets a 60 GB hard limit; exceeding it causes the actor to spill GPU objects to CPU DRAM with a 10% throughput penalty for subsequent accesses. The memory manager tracks allocation via `ray._private.gpu_memory` and logs warnings when utilization exceeds 80% of quota. On multi-GPU nodes, Ray automatically assigns GPU IDs using `CUDA_VISIBLE_DEVICES` isolation, preventing GPU memory conflicts between actors on the same node.
| Ray GPU Feature | Ray 2.25 | Ray 2.30 | Ray 2.33+ |
|---|---|---|---|
| GPU Direct Object Store | Not supported | GPUDirect RDMA (beta) | GPUDirect RDMA (GA) |
| GPU Memory Quotas | Not supported | RAY_gpu_memory_manager=1 | Default (70% quota) |
| GPU Memory Spill to CPU | No (OOM on overflow) | Yes (10% perf penalty) | Yes (configurable penalty) |
| CUDA_VISIBLE_DEVICES Isolation | Manual | Automatic (actor-level) | Automatic (task-level) |
| GPU Utilization Monitoring | ray logs only | ray metrics + Prometheus | Per-GPU + per-process |
| NVIDIA NCCL Support | Limited (TCP only) | NCCL fast socket | NCCL IB + GDR |
RAY SERVE: GPU-AWARE INFERENCE DEPLOYMENT WITH AUTOSCALING
Ray Serve deploys GPU models as deployment replicas with `@serve.deployment(num_gpus=1)` annotations. Each replica requests one GPU via Ray's resource scheduler, which uses a bin-packing strategy to maximize GPU utilization: it packs multiple replicas onto a single GPU only if `num_gpus=0.5` or fractional values are specified. For Llama 70B inference, the typical deployment uses `num_gpus=4` with `ray_actor_options={"runtime_env": {"env_vars": {"CUDA_VISIBLE_DEVICES": "0,1,2,3"}}}` to pin the replica to 4 specific GPUs for tensor parallelism. Ray Serve's autoscaler monitors queue depth and GPU utilization: when GPU utilization exceeds 85% for 60 seconds, the autoscaler adds a new replica (which may require provisioning a new node from the cluster autoscaler). The min_replicas parameter for critical inference should be set such that peak load minus min_replicas * per_replica_throughput = 30% headroom.
The Ray Serve autoscaler integrates with Kubernetes via KubeRay, which manages RayCluster custom resources. KubeRay's `rayClusterSpec.workerGroupSpecs[].numOfHosts` maps to GPU node pools: a typical configuration defines a `gpu-h100` worker group with `maxReplicas=32` and `acceleratorType="nvidia.com/gpu"` for GPU scheduling. KubeRay 1.2+ adds `scaleStrategy: "TRIGGER"` for reactive scaling (scale up on queue depth, scale down on idle) versus `"PREFETCH"` for proactive scaling (pre-warm replicas based on predicted traffic). For production inference, the TRIGGER mode with a 120-second cooldown prevents replica flapping under variable load. Ray Serve also supports request batching via `@serve.deployment(max_batching_size=32, batch_wait_timeout_s=0.1)`, accumulating up to 32 requests or waiting 100ms before dispatching to the GPU model.
RAY TRAIN: ELASTIC DISTRIBUTED TRAINING WITH FAULT TOLERANCE
Ray Train provides framework-agnostic distributed training with `TorchTrainer`, `TransformersTrainer`, and the newer `DataParallelTrainer` for custom training loops. The `TorchTrainer` wraps PyTorch DDP with automatic NCCL setup: `trainer = TorchTrainer(train_func, scaling_config=ScalingConfig(num_workers=8, use_gpu=True, resources_per_worker={"GPU": 1}))`. The `ScalingConfig` defines the degree of data parallelism (number of workers) and resource allocation per worker. Ray Train handles NCCL rendezvous and world group setup automatically via a Ray actor group, replacing the manual `torchrun` or `torch.distributed.init_process_group` boilerplate. For Llama 70B training on 8x H100, Ray Train sets up NCCL with NVIDIA's NCCL fast socket transport and GPUDirect for cross-node all-reduce.
Ray Train's fault tolerance is its differentiator from native PyTorch DDP. When a GPU worker fails during training (e.g., NCCL timeout, GPU memory OOM, node failure), Ray Train detects the failure via the Ray GCS (Global Control Store) heartbeat mechanism and triggers a coordinated recovery. The latest checkpoint is loaded from the shared filesystem (NFS or S3), the failed worker is replaced with a newly allocated GPU, and NCCL communicators are rebuilt with the new world size. The recovery time depends on checkpoint frequency: with 5-minute checkpoint intervals, a single-GPU failure on an 8-GPU training run incurs 6-8 minutes of lost training time (5-min checkpoint gap + 1-min checkpoint reload + 2-min NCCL rebuild). Ray Train's elastic scaling extends this: `trainer = TorchTrainer(..., scaling_config=ScalingConfig(num_workers=8, max_num_workers=16))` allows the trainer to dynamically add workers during training (e.g., when spot instance pricing drops), increasing training throughput without restarting.
| Ray Train Feature | PyTorch DDP Equivalent | Benefit |
|---|---|---|
| Auto NCCL Rendezvous | torchrun / init_process_group | No manual MASTER_ADDR setup |
| Fault Tolerance | Manual retry loop | 6-8 min recovery vs 30+ min manual |
| Elastic Scaling | Not supported | Scale workers mid-training |
| Mixed Precision (FP8) | torch.cuda.amp + TE | Integrated with Ray Train config |
| Checkpoint Restart | torch.save/load manual | Auto restore to worker count |
| GPU Memory Reporting | nvidia-smi polling | Per-worker + per-step logs |
RAY CLUSTER CONFIGURATION AND GPU AUTOSCALING ON KUBERNETES
Production Ray GPU cluster configuration requires tuning five parameters. `object_store_memory` (default: 4 GB) should be set to 20-30% of total GPU memory for GPU workloads: `ray.init(object_store_memory=20*1024**3)` for an H100 node. `num_cpus_per_worker` and `num_gpus_per_worker` must match Kubernetes resource limits. `RAY_gcs_server_rpc_timeout_seconds` should be increased from default 60s to 300s for GPU training runs: NCCL barriers can block GCS heartbeats for extended periods. `RAY_driver_worker_heartbeat_timeout` set to 600s prevents false-positive worker eviction during long GPU kernel execution. `RAY_max_pending_launch_requests` set to 128 controls how many GPU workers can be pending allocation simultaneously, preventing overwhelming the Kubernetes scheduler during rapid scale-up events.
Kubernetes GPU autoscaling with Ray uses the Kubernetes Cluster Autoscaler with GPU node group configuration. The autoscaler monitors pending Ray pod requests with GPU resource requirements and scales a `gpu-h100` node group when pending GPU requests exceed a threshold (typically 8 GPUs = 1 node). The `cluster-autoscaler.kubernetes.io/safe-to-evict: "false"` annotation on Ray pods prevents premature eviction of GPU training workers during node scale-down. A GPU node provisioning delay of 3-7 minutes (cloud provider-dependent) means the autoscaler should provision 2-3 nodes ahead of demand by setting `RAY_preprovision_nodes=2`. On ClusterBid, Ray GPU nodes are provisioned with GPU-optimized instances: 8x H100 nodes at $20/hr for training, single H100 nodes at $2.50/hr for Ray Serve inference, with autoscaling policies that target 75% GPU utilization steady state.
RAY DATA AND GPU-ACCELERATED DATA PIPELINES
Ray Data (formerly Ray Datasets) provides GPU-accelerated data preprocessing for training pipelines. The `ds.map_batches(train_func, batch_size=32, num_gpus=1)` API processes batches on GPU, using the same `num_gpus` resource annotation as Ray Train and Ray Serve. For vision transformer training, `ds.map_batches(augment_images, batch_size=128, num_gpus=2, compute=ray.data.ActorPoolStrategy(min_size=2, max_size=8))` distributes image augmentation across 2-8 GPU workers processing 128-image batches each. The GPU-accelerated data pipeline achieves 3-5x faster preprocessing than CPU-only: JPEG decoding with nvJPEG on GPU processes 1,200 images/second per GPU versus 350 images/second on a 16-core CPU node. Ray Data pipelines automatically fuse preprocessing operators when possible, reducing GPU kernel launch overhead.
The `ray.data.DatasetWriter` supports writing to GPU memory directly using Arrow format with CUDA integration. For training datasets that do not fit in GPU memory, Ray Data supports streaming reads with prefetching: `ds = ray.data.read_parquet("s3://dataset/", shuffle="files", prefetch_batches=4)`. The prefetch buffer (4 batches per GPU worker) consumes GPU memory for intermediate data, which must be accounted for in GPU memory budgets. The `ephemeral_oom_retry` parameter in `map_batches` retries batches that fail due to out-of-memory errors with reduced batch size, preventing pipeline crashes from GPU memory pressure. Ray Data's integration with Ray Train means the data pipeline and training loop share the same GPU cluster, reducing data transfer latency to near-zero by keeping preprocessed tensors in GPU memory across pipeline stages.
