Skip to content

Latest commit

 

History

83 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Ray Hive

Ray Hive is a small control layer for running a heterogeneous Ray cluster on k3s and placing vLLM model replicas according to live GPU memory.

The cluster remains a normal Ray cluster: CPU workers can run general distributed tasks while GPU workers host model deployments through Ray Serve.

Contents

Architecture

Clients
  └─ Ray head
      ├─ CPU workers → general Ray tasks
      └─ GPU workers → vLLM replicas
           ├─ VRAM monitor DaemonSet
           ├─ shared GPU registry
           └─ Ray Serve model routers

KubeRay manages the head and worker pods. Ray handles task scheduling and Serve deployments; Ray Hive adds model planning, per-GPU placement, memory reservations, and routing.

Capabilities

  • Throughput-first planning — max practical max_num_seqs / batched tokens within a VRAM budget (gpu_budget_frac 0.95). With max_input_prompt_length="auto" + fixed max_num_seqs, grow text context (floor 256) instead. At max_num_seqs=1, auto may shrink below 256 so weights + speculative draft still fit. Hybrid GDN/Mamba models cap auto input at 32768 so leftover VRAM is not spent on an unbounded attention window.
  • Heterogeneous replicas — per-GPU plans so different cards each contribute what they can.
  • Least-loaded routing — relative to each replica’s planned capacity.
  • Live VRAM scheduling — registry + reservations; see examples/4_test_allocation_policies.py.
  • GPU sharing — co-locate when needed; intentional share via same gpu= pin (examples/6_shared_gpu.py).
  • Flexible placement — pin, N replicas, or replicas=-1 (examples/1_test_model_configs.py).
  • Flagship deploy reference — Qwen3.8-27B + Nemotron 3.5 Lightning variants (examples/15_flagship_models.py).
  • Same-node TP — auto escalate or pin a GPU list (examples/5_tensor_parallel.py).
  • Custom attention — subclass using HF config fields (examples/3_custom_attention.py).
  • Multimodal generate — image / video / audio (examples/8–9, 12–14).
  • Embeddings — runner="pooling" (examples/10–11).
  • Sleep / idle — level-1 sleep then optional self-destroy (examples/7_sleep_idle_timeout.py).
  • OpenAI HTTP — per-model routes + cluster /v1 gateway (examples/2_test_inference.py).
  • General Ray tasks — same cluster for CPU work.

Out of scope: precomputed multimodal embeds as input (enable_mm_embeds).

Quick start

Dry-run packing against live GPUs (same plan deploy would use):

from ray_hive import RayHive

hive = RayHive(address="ray://YOUR_RAY_HEAD_IP:10001")
plan = hive.estimate_vram(
    "Qwen/Qwen3-0.6B-FP8",
    max_input_prompt_length=1024,
    max_output_prompt_length=2048,
    replicas=1,
    vllm_kwargs={
        "trust_remote_code": True,
        "reasoning_parser": "qwen3",
        "default_chat_template_kwargs": {"enable_thinking": False},
    },
)
Deployment Plan: Qwen/Qwen3-0.6B-FP8
  Replica                  replica-0
  GPU(s)                   host-a:gpu0
  tensor_parallel_size     1
  max_input_prompt_length  1024
  max_output_prompt_length 2048
  max_model_len            3072
  max_num_seqs             48
  max_num_batched_tokens   8192
  gpu_memory_utilization   0.850
  Weights                  0.60 GiB
  KV cache                 4.20 GiB
  Activations              0.40 GiB
  Overhead                 0.30 GiB
  Total (per GPU)          5.50 GiB

Deploy until ready (router up + warmed):

status = hive.deploy_model(
    model_id="qwen",
    model_name="Qwen/Qwen3-0.6B-FP8",
    max_input_prompt_length=1024,
    max_output_prompt_length=2048,
    replicas=1,
    vllm_kwargs={
        "trust_remote_code": True,
        "reasoning_parser": "qwen3",
        "default_chat_template_kwargs": {"enable_thinking": False},
    },
)
{"model_id": "qwen", "status": "ready", "route": "/qwen", "openai_v1": "/v1", "replicas": {...}}

Generate text (batch and stream helpers available):

from ray_hive.inference import inference, inference_batch, inference_stream

answer = inference("Explain Ray in one sentence.", model_id="qwen")
answers = inference_batch(["Prompt one", "Prompt two"], model_id="qwen")
for delta in inference_stream("Count to three.", model_id="qwen"):
    print(delta, end="")
Ray is a distributed computing framework for scaling Python workloads.

Structured output:

from pydantic import BaseModel

class Answer(BaseModel):
    summary: str
    confidence: float

result = inference("Summarize Ray.", model_id="qwen", structured_output=Answer)

Inside an asyncio loop use a_inference / a_inference_batch instead of the sync helpers.

Placement

  • gpu="host:gpu0" — one replica on one GPU.
  • gpu=[a,b,...] + replicas=len(list) — N single-GPU pins (TP=1 each).
  • gpu=[a,b,...] + replicas=1 — one same-node TP group.
  • gpu=None — auto place (single GPU via allocation_cls, default RayPerformanceAllocator; else same-node TP).
  • replicas=-1 — every eligible GPU / TP group.

Auto input length

By default context lengths are a fixed contract and omitted max_num_seqs is packed to fill VRAM.

max_input_prompt_length="auto" does the inverse for text input only: you must pass max_num_seqs in vllm_kwargs, output length stays fixed, and the planner grows text input from a floor of 256 until the VRAM budget is filled (capped by HF max_position_embeddings / model_max_length when present). When max_num_seqs=1 and even 256 does not leave KV room after weights + speculative draft, auto shrinks toward 1 token. Multimodal placeholder tokens are still added on top of the chosen text length, so max_model_len always covers MM + output.

speculative_config is part of the memory plan: in-checkpoint MTP reserves an extra compute-dtype lm_head (and num_nextn_predict_layers when set); a separate draft model adds that checkpoint’s weights and KV.

hive.estimate_vram(
    "Qwen/Qwen3-0.6B-FP8",
    max_input_prompt_length="auto",
    max_output_prompt_length=512,
    replicas=1,
    vllm_kwargs={
        "max_num_seqs": 8,
        "trust_remote_code": True,
        "reasoning_parser": "qwen3",
        "default_chat_template_kwargs": {"enable_thinking": False},
    },
)
  max_input_prompt_length  8192 (auto)
  max_output_prompt_length 512
  max_model_len            8704
  max_num_seqs             8

Passing "auto" without max_num_seqs raises ConfigError.

Pin one GPU:

hive.deploy_model(
    model_id="qwen-pin",
    model_name="Qwen/Qwen3-0.6B-FP8",
    max_input_prompt_length=512,
    max_output_prompt_length=512,
    replicas=1,
    gpu="ergos-06-nv:gpu0",
    vllm_kwargs={"trust_remote_code": True, "reasoning_parser": "qwen3",
                 "default_chat_template_kwargs": {"enable_thinking": False}},
)

All eligible GPUs:

hive.deploy_model(
    model_id="qwen-all",
    model_name="Qwen/Qwen3-0.6B-FP8",
    max_input_prompt_length=512,
    max_output_prompt_length=512,
    replicas=-1,
    vllm_kwargs={"trust_remote_code": True, "reasoning_parser": "qwen3",
                 "default_chat_template_kwargs": {"enable_thinking": False}},
)

Same-node TP=2 pin:

hive.deploy_model(
    model_id="qwen-tp",
    model_name="Qwen/Qwen3-8B-FP8",
    max_input_prompt_length=512,
    max_output_prompt_length=512,
    replicas=1,
    gpu=["ergos-02-nv:gpu1", "ergos-02-nv:gpu2"],
    vllm_kwargs={"trust_remote_code": True, "reasoning_parser": "qwen3",
                 "default_chat_template_kwargs": {"enable_thinking": False}},
)

Live registry snapshot: hive.get_vram_state(). Policies: Allocation policies.

Allocation policies

When gpu= is unset, allocation_cls picks GPUs (ray_hive.core.ray_gpu_alloc). Pins skip policy ranking but still check VRAM fit and arch taints. Default: RayPerformanceAllocator.

  • RayPerformanceAllocator — rank by compute (SM count, light bandwidth tie-break); top-N.
  • RayConserveTdpAllocator — prefer lower approx TDP; SM count tie-break.
  • RayTensorParallelAllocator — same-node packs of size tensor_parallel_size; auto when single-GPU fails.

Arch taint: native FP8 (HF / vLLM dtype / kv / quantization mentioning fp8 / float8 / float-quantized) needs compute capability ≥ 8.9 (Ada+); Ampere dropped automatically.

TP=1 policies prefer unshared GPUs first, then co-locate. Alive Ray nodes only. See examples/4_test_allocation_policies.py and examples/5_tensor_parallel.py.

Custom attention

VRAM planning sizes the KV cache through an attention-specs class. By default that is BaseAttentionSpecs (or MultimodalAttentionSpecs for multimodal models). Pass a subclass via attention_cls on estimate_vram / deploy_model.

The class sees the same HF config fields the planner already loads. Override properties or methods that use those fields when the default transformer KV math does not match the model (e.g. heads, sliding windows, extra tokens).

from ray_hive.core.model_specs import BaseAttentionSpecs

class MyAttention(BaseAttentionSpecs):
    @property
    def kv_heads(self) -> int:
        return self.hf_params["num_key_value_heads"]  # from HF config

hive.deploy_model(
    ...,
    attention_cls=MyAttention,
)

See examples/3_custom_attention.py (and examples/14_gemma4_stress.py for multimodal).

Multimodal and context

Text-only models — pass a string (or text chat messages). No limit_mm_per_prompt.

MM models — planner uses limit_mm_per_prompt in vllm_kwargs to size worst-case image / video / audio placeholders. If omitted on an MM HF config, defaults are derived from vision_config / audio_config (typically image: 1 and/or audio: 1). Set counts explicitly to enable or disable modalities.

Token budget — max_input_prompt_length is the text side only. Effective input ≈ text + MM placeholders; max_model_len ≈ effective_input + max_output_prompt_length (use output 0 for pooling). If that cannot cover placeholders + output, planning raises MmContextError — raise max_input_prompt_length or lower limit_mm_per_prompt. With max_input_prompt_length="auto" (requires max_num_seqs), text grows from 256 (or shrinks below that at max_num_seqs=1) while MM placeholders stay reserved inside max_model_len.

Requests vs planning — on an MM deploy you can still send text-only strings. That does not shrink the VRAM plan. For text-only planning on an MM checkpoint, zero unused modalities ({"image": 0, "video": 0, "audio": 0}).

Enable image chat:

from ray_hive.core.ray_utils import file_to_data_url
from ray_hive.inference import inference

hive.deploy_model(
    model_id="vl",
    model_name="Qwen/Qwen2.5-VL-3B-Instruct",
    max_input_prompt_length=2048,   # text tokens you expect
    max_output_prompt_length=256,
    replicas=1,
    vllm_kwargs={
        "trust_remote_code": True,
        "limit_mm_per_prompt": {"image": 1},
    },
)

messages = [{
    "role": "user",
    "content": [
        {"type": "image_url", "image_url": {"url": file_to_data_url("photo.png")}},
        {"type": "text", "text": "Describe this image."},
    ],
}]
answer = inference(messages, model_id="vl", max_tokens=64)
# text-only request on the same deploy still works:
inference("Say hello.", model_id="vl", max_tokens=32)
# estimate_vram on MM models also shows:
  mm_tokens_per_prompt     1280

Content part types: text, image_url, video_url / video, audio_url / input_audio. Fixtures: examples/media/. Demos: examples/8_multimodal_vision.py, 9_multimodal_audio.py, 12_multimodal_video.py, 13_gemma4_multimodal.py, 14_gemma4_stress.py.

Embeddings

Set runner="pooling" (legacy task="embed" still works). Use max_output_prompt_length=0. inference returns a vector (or list of vectors for batch).

hive.deploy_model(
    model_id="embed",
    model_name="BAAI/bge-small-en-v1.5",
    max_input_prompt_length=512,
    max_output_prompt_length=0,
    replicas=1,
    vllm_kwargs={"runner": "pooling", "trust_remote_code": True},
)
vec = inference("hello", model_id="embed")
[-0.012, 0.034, ...]   # length = model hidden size
curl http://YOUR_RAY_HEAD_IP:8000/embed/v1/embeddings \
  -H "Content-Type: application/json" \
  -d '{"model":"embed","input":["hello world"]}'

See examples/10_text_embeddings.py and examples/11_multimodal_embeddings.py.

Lifecycle

Sleep after quiet, then optional full destroy (idle_timeout must be greater than sleep_timeout when both are set). While sleeping, the running footprint stays reserved so other deploys cannot take the space a wake needs. That reservation is one copy of the weights plus the planned KV cache. Pass sleep_peak_factor in vllm_kwargs only if you need an extra weight-sized hold:

hive.deploy_model(
    model_id="qwen",
    model_name="Qwen/Qwen3-0.6B-FP8",
    max_input_prompt_length=512,
    max_output_prompt_length=256,
    replicas=1,
    sleep_timeout=60,
    idle_timeout=300,
    vllm_kwargs={
        "trust_remote_code": True,
        "reasoning_parser": "qwen3",
        "default_chat_template_kwargs": {"enable_thinking": False},
        # "sleep_peak_factor": 1.0,  # optional extra weight-sized hold
    },
)
hive.shutdown("qwen")   # one model
hive.shutdown()         # all models
from ray_hive import kill_gpu_registry
kill_gpu_registry()     # force registry rebuild (DaemonSet re-registers)

See examples/7_sleep_idle_timeout.py and examples/0_shutdown_models.py.

OpenAI HTTP

Per-model route and cluster-wide /v1 gateway (port 8000). Gateway starts on deploy_model; for an already-running cluster call hive.ensure_openai_api().

curl http://YOUR_RAY_HEAD_IP:8000/qwen/v1/chat/completions \
  -H "Content-Type: application/json" \
  -d '{"model":"qwen","messages":[{"role":"user","content":"Explain Ray briefly."}]}'

# Open WebUI: API URL = http://YOUR_RAY_HEAD_IP:8000/v1
curl http://YOUR_RAY_HEAD_IP:8000/v1/models
curl http://YOUR_RAY_HEAD_IP:8000/v1/chat/completions \
  -H "Content-Type: application/json" \
  -d '{"model":"qwen","messages":[{"role":"user","content":"Explain Ray briefly."}]}'

Streaming, structured JSON, and LangChain: examples/2_test_inference.py.

General Ray CPU Tasks

import ray

ray.init(address="ray://YOUR_RAY_HEAD_IP:10001")

@ray.remote
def square(value):
    return value * value

results = ray.get([square.remote(value) for value in range(10)])

Running tests

pip install -e ".[dev]"   # or: pip install -r requirements-dev.txt

# Parallel unit suite (default; no Ray / no LLM — RAY_ADDRESS not required)
pytest -n auto -m "not ray and not live"

ray / live require RAY_ADDRESS. Without it those tests skip. Set it in the shell before running (same value as examples/.env):

# Linux / macOS
export RAY_ADDRESS=ray://YOUR_RAY_HEAD_IP:10001

# PowerShell
$env:RAY_ADDRESS = "ray://YOUR_RAY_HEAD_IP:10001"

Then (single process — do not use pytest -n here; live GPU claims must stay process-local):

pytest -m ray    # cluster smoke, no LLM
pytest -m live   # Deploy A–H + auto-input + sleep-hold; needs GPUs + model download
# Focused:
#   pytest -m live -k "auto_input or sleep_vram_hold"
# or both:
pytest -m "ray or live"

Markers: ray (cluster, no LLM), live (small-model cycles), nightly (stress / resilience).

Repository

  • ray_hive/ — planner, registry, deployment service, router, client API
  • manifests/ — KubeRay cluster, worker image, VRAM monitor
  • examples/ — deploy/inference experiments (.env / .env.example)
  • examples/requirements.txt — example-only deps
  • examples/media/ — multimodal fixtures
  • tests/ — pytest unit (unit/), Ray smoke (ray_smoke/), live cycles (live/)
  • basic_ray_tests/ — manual cluster / resource scripts

Related: rayify, a tool for converting scripts into Ray jobs.

About

Ray Hive is a distributed LLM serving SDK for a Ray cluster on k3s. It deploys vLLM models across heterogeneous GPUs with VRAM-aware scheduling, then exposes a simple Python API for inference. Included: Infrastructure as Code with Kubernetes manifests, Helm charts, and deployment scripts.

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages