Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Important
Ez a funkció nyilvános előzetes verzióban van.
Ez a példa offline LLM batch inferenciát futtat Ray Data és vLLM segítségével 8 A10 csomóponton keresztül. Egy bootstrap szkript indít egy Ray klasztert a csomópontok között, majd az illesztőprogram a Ray Data LLM API-ját (ray.data.llm) használja, hogy minden csomóponton egy vLLM replikát indítson, és egy adathalmazt adjon át rajtuk, a generált szöveget pedig Parquet néven írja egy Unity Catalog kötetre.
Egy nyilvános modellt (Qwen2.5-7B-Instruct) használ, így változtatás nélkül futtatható Hugging Face-token nélkül.
A számítási feladat a következőket végzi el:
- Feltölti a helyi projektet a következővel
code_source: snapshot: . - Indít egy Ray fejet a 0-es csomóponton, csatlakozik 7 munkacsomóponthoz, majd futtatja a batch inference drivert.
- A
ray.data.llmhasználatával csomópontonként egy-egy vLLM replika fut, és a promptok feldolgozása párhuzamosan történik. - Parquetként írja a parancssorokat és a létrehozott kimeneteket egy Unity Catalog-kötetbe.
Prerequisites
- A parancssori
airfelület telepítve és hitelesítve. Lásd : Az AI-futtatókörnyezeti parancssori felület telepítése. - Egy Unity Catalog-kötet, amelybe írhat. Az elérési útját az alábbi munkaterhelés YAML-fájljában állíthatja be.
Projektelrendezés
Hozzon létre egy könyvtárat a következő fájlokkal.
ray_batch_inference/
├── train.yaml # air workload config (inline dependencies + Ray bootstrap)
└── batch_inference.py # Ray Data + vLLM batch inference driver
1. lépés: A számítási feladat YAML-jének írása
train.yaml 8 GPU_1xA10 csomópontot kér. A függőségeket közvetlenül a environment alatt adják meg (a version kliensképpel), és a command Ray-klasztert indít a csomópontokon, majd futtatja a drivert, így a munkaterhelésnek nincs szüksége külön függőségi fájlra vagy indítószkriptre.
A vLLM nincs benne a base image-ben, ezért inline van telepítve, együtt azzal a három rögzített függőséggel, amelyekre a GPU-csomópontoknak szükségük van: hf_transfer (a base image lehetővé teszi a gyors Hugging Face-letöltéseket, és ezt a csomagot várja el), egy újabb fsspec (a base image egy régi verziót tartalmaz, amely tönkreteszi a letöltéseket), valamint egy rögzített opencv-python-headless-verzió (a vLLM behúzza az OpenCV-t, amelynek alapértelmezett wheel csomagja meghiúsítja az OpenSSL FIPS öntesztjét a GPU-csomópontokon).
Állítsa a(z) OUTPUT_PATH elemet egy olyan Unity Catalog-kötetre, amelyre írhat. Állítsuk NUM_GPUS ugyanarra az értékre, mint num_accelerators.
experiment_name: air-ray-batch-inference
environment:
version: '5'
dependencies:
- ray[data]==2.56.1
- vllm
- datasets>=3.0
- huggingface_hub>=0.34
# The base image sets HF_HUB_ENABLE_HF_TRANSFER=1; install the package it expects
# so model and dataset downloads don't error out.
- hf_transfer
# The base image ships fsspec 2023.5.0, which is too old for modern
# huggingface_hub and breaks dataset/model downloads. Pin a newer fsspec.
- fsspec>=2024.6.1
# vLLM pulls in opencv; its default wheel crashes the OpenSSL FIPS self-test
# on the GPU nodes. This pinned headless build avoids the crash.
- opencv-python-headless==4.12.0.88
# 8 A10 nodes, one GPU each. Ray Data runs one vLLM replica per node.
compute:
num_accelerators: 8
accelerator_type: GPU_1xA10
code_source:
type: snapshot
snapshot:
root_path: .
command: |
cd $CODE_SOURCE_PATH
RAY_HEAD_PORT=6379
GPUS_PER_NODE=${LOCAL_WORLD_SIZE:-1}
if [ "${NODE_RANK:-0}" = "0" ]; then
echo "NODE_RANK=0: starting Ray head with $GPUS_PER_NODE GPU(s)..."
ray start --head --port=$RAY_HEAD_PORT --num-gpus="$GPUS_PER_NODE" --dashboard-host=0.0.0.0
python batch_inference.py
ray stop
else
echo "NODE_RANK=$NODE_RANK: connecting to Ray head at $MASTER_ADDR:$RAY_HEAD_PORT..."
for i in $(seq 1 12); do
if ray start --address="$MASTER_ADDR:$RAY_HEAD_PORT" --num-gpus="$GPUS_PER_NODE" --block 2>/dev/null; then
break
fi
echo "Attempt $i failed, retrying in 5s..."
sleep 5
done
fi
max_retries: 0
timeout_minutes: 60
env_variables:
NCCL_SOCKET_IFNAME: eth0
# Unity Catalog volume where results land as Parquet. Replace with your volume.
OUTPUT_PATH: /Volumes/main/default/air_examples/ray_batch_inference
NUM_GPUS: '8' # must match num_accelerators
Az inline command elindít egy Ray-headet a 0. csomóponton lévő GPU-val, futtatja a drivert a python batch_inference.py használatával, majd leállítja a klasztert. A feldolgozócsomópontok a fejcsomóponthoz a(z) MASTER_ADDR és NODE_RANK használatával csatlakoznak, amelyeket a platform automatikusan állít be.
2. lépés: A köteg következtetési illesztőprogramjának meghatározása
batch_inference.py létrehoz egy Ray-adatkészletet promptokból, konfigurál egy vLLM-processzort a(z) ray.data.llm használatával, és kiírja az eredményeket.
concurrency a Ray Data párhuzamosan futó vLLM-replikák száma. Az illesztőprogram megvárja, hogy minden csomópont csatlakozzon, mielőtt elolvassa a GPU számát, így minden csomópontot használnak:
import os
import time
import ray
from ray.data.llm import build_processor, vLLMEngineProcessorConfig
ray.init(address="auto")
num_gpus = int(os.environ["NUM_GPUS"])
for _ in range(60):
if int(ray.cluster_resources().get("GPU", 0)) >= num_gpus:
break
time.sleep(5)
total_gpus = int(ray.cluster_resources().get("GPU", 0))
if total_gpus < num_gpus:
raise SystemExit(f"Expected {num_gpus} GPU(s) but Ray only sees {total_gpus}.")
config = vLLMEngineProcessorConfig(
model_source="Qwen/Qwen2.5-7B-Instruct",
engine_kwargs={"max_model_len": 4096, "tensor_parallel_size": 1},
concurrency=total_gpus, # one vLLM replica per GPU in the cluster
batch_size=64,
)
processor = build_processor(
config,
preprocess=lambda row: dict(
messages=[{"role": "user", "content": row["instruction"]}],
sampling_params=dict(max_tokens=256, temperature=0.7),
),
postprocess=lambda row: dict(instruction=row["instruction"], output=row["generated_text"]),
)
out = processor(ds) # ds is a Ray Dataset with an "instruction" column
out.write_parquet(OUTPUT_PATH)
preprocess az egyes bemeneti sorokat csevegési kéréssé alakítja, és postprocess megőrzi az oszlopokat. A Ray Data hozzáad egy oszlopot generated_text a modell kimenetével. A teljes szkript a lap végén található teljes illesztőprogram-szkriptben található.
Nagyobb modellek esetén állítsa a tensor_parallel_size értékét úgy, hogy egy replika több GPU között legyen felosztva, és ossza el a total_gpus értékét ezzel az értékkel, hogy a replikák továbbra is teljesen kihasználják a fürtöt, például concurrency=total_gpus // 2 és tensor_parallel_size=2 használatával.
3. lépés: A futtatás beküldése
air run -f train.yaml --dry-run
air run -f train.yaml --watch
4. lépés: A futtatás vizsgálata
air get run <run-id>
air logs <run-id>
A naplók a vLLM-motor parancssorát és generációjának átviteli sebességét jelenítik meg a köteg futtatásakor, majd a kimenet írásakor egy Wrote <n> rows sort.
Ahol az eredmények megjelennek
Az illesztőprogram egy Parquet-adatkészletet ír a OUTPUT_PATH kötetbe, egy instruction oszloppal és egy output oszloppal. Olvassa be újra a Sparkkal vagy a pandasszal, például spark.read.parquet(OUTPUT_PATH).
Teljes illesztőprogram-szkript
A teljes batch_inference.py másolás-beillesztéshez:
#!/usr/bin/env python3
"""Offline batch inference with Ray Data + vLLM across 8 A10 nodes.
The workload `command` starts a Ray head on node 0 and joins 7 worker nodes, each
contributing 1 GPU. Ray Data's LLM API (`ray.data.llm`) launches one vLLM replica
per GPU and streams a dataset of prompts through them, then writes the generated text
to a Unity Catalog volume as Parquet.
Uses a public model (no Hugging Face token required) so the example runs as-is.
"""
import os
import time
import ray
from datasets import load_dataset
from ray.data.llm import build_processor, vLLMEngineProcessorConfig
MODEL_SOURCE = "Qwen/Qwen2.5-7B-Instruct"
NUM_PROMPTS = 1000
# Unity Catalog volume path where results land as Parquet. Set this in train.yaml.
OUTPUT_PATH = os.environ.get("OUTPUT_PATH", "/Volumes/main/default/air_examples/ray_batch_inference")
def build_prompts():
"""Build a Ray Dataset of prompts from a public instruction dataset."""
raw = load_dataset("tatsu-lab/alpaca", split=f"train[:{NUM_PROMPTS}]")
items = []
for row in raw:
instruction = row["instruction"]
if row.get("input"):
instruction = f"{instruction}\n\n{row['input']}"
items.append({"instruction": instruction})
return ray.data.from_items(items)
def main():
ray.init(address="auto")
num_gpus = int(os.environ["NUM_GPUS"])
for _ in range(60):
if int(ray.cluster_resources().get("GPU", 0)) >= num_gpus:
break
time.sleep(5)
total_gpus = int(ray.cluster_resources().get("GPU", 0))
if total_gpus < num_gpus:
raise SystemExit(
f"Expected {num_gpus} GPU(s) but Ray only sees {total_gpus}; "
"check GPU discovery / node join on all nodes."
)
print(f"Ray cluster ready: {total_gpus} GPU(s)", flush=True)
ds = build_prompts()
# vLLM engine config. concurrency = number of replicas Ray Data runs in parallel;
# one per GPU in the cluster here. engine_kwargs are passed through to the vLLM engine.
config = vLLMEngineProcessorConfig(
model_source=MODEL_SOURCE,
engine_kwargs={
"max_model_len": 4096,
"tensor_parallel_size": 1,
"enable_chunked_prefill": True,
},
concurrency=total_gpus,
batch_size=64,
)
# preprocess maps each input row to a chat request; postprocess keeps the columns
# we want to persist. ray.data.llm adds a `generated_text` column.
processor = build_processor(
config,
preprocess=lambda row: dict(
messages=[
{"role": "system", "content": "You are a helpful assistant."},
{"role": "user", "content": row["instruction"]},
],
sampling_params=dict(max_tokens=256, temperature=0.7),
),
postprocess=lambda row: dict(
instruction=row["instruction"],
output=row["generated_text"],
),
)
# materialize once so the write and the sample print don't re-run inference.
out = processor(ds).materialize()
out.write_parquet(OUTPUT_PATH)
print(f"Wrote {out.count()} rows to {OUTPUT_PATH}", flush=True)
for row in out.take(2):
print("INSTRUCTION:", row["instruction"][:120], flush=True)
print("OUTPUT:", row["output"][:200], flush=True)
ray.shutdown()
if __name__ == "__main__":
main()