Qwen2.5-32B Batch-Inferenz mit Ray Data und vLLM

Verwenden Sie Qwen2.5-32B-Instruct, um 16.000 mehrsprachige Sprachassistenten-Äußerungen auf einer angehängten 8xH100 KI-Laufzeit zu klassifizieren. Dieses Notizbuch zeigt, wie man:

  • Baue ein ausgewogenes mehrsprachiges Ray-Dataset aus MASSIVE 1.1.
  • Führe auf jeder verfügbaren GPU eine persistente vLLM-Modellreplik aus.
  • Überwachen Sie die Arbeitsbelastung mit dem Ray-Dashboard und den MLflow-Systemmetriken.
  • Speichern Sie vollständige Vorhersageergebnisse als Parquet in einem Volume von Unity Catalog.

Note

Dieses Beispiel benötigt die Databricks KI-Umgebung Version 5 oder höher.

Verbindung zu Serverless GPU-Compute herstellen

  1. Im Notebook-Compute-Selektor wählen Sie Serverless GPU aus.
  2. Im Panel Umgebung wählen Sie den 8xH100-Beschleuniger und die KI-v5-Umgebung aus.
  3. Klicken Sie auf Anwenden und bestätigen Sie dann die Umgebung.

Das Qwen-Modell ist öffentlich und erfordert keine Hugging Face-Authentifizierung. Das Notizbuch lädt MASSIVE 1.1 aus dem öffentlichen Amazon-Archiv herunter.

Importieren von Bibliotheken

AI v5 enthält die Ray-, vLLM-, Hugging Face Datasets-, Transformers-, PyTorch- und MLflow-Pakete, die in diesem Notebook verwendet werden, sodass keine Paketinstallation erforderlich ist.

import json
import re
import time
from pathlib import Path

import mlflow
import pandas as pd
from datasets import DownloadConfig, DownloadManager, concatenate_datasets, load_dataset
from datasets.utils.logging import disable_progress_bar
from pyspark.sql import functions as F
from vllm import LLM, SamplingParams
from vllm.sampling_params import StructuredOutputsParams

Konfigurieren Sie die Arbeitslast

Stellen Sie das Modell, die Orte, die Stichprobengröße und die Inferenzparameter fest.

MODEL_NAME = "Qwen/Qwen2.5-32B-Instruct"
DATASET_NAME = "AmazonScience/massive"
MASSIVE_ARCHIVE_URL = "https://amazon-massive-nlu-dataset.s3.amazonaws.com/amazon-massive-dataset-1.1.tar.gz"
LOCALES = ["en-US", "es-ES", "de-DE", "ar-SA", "hi-IN", "ja-JP", "sw-KE", "zh-CN"]
ROWS_PER_LOCALE = 2_000
BATCH_SIZE = 64
MAX_MODEL_LEN = 512
MAX_OUTPUT_TOKENS = 8
SEED = 42

Konfiguration des Unity-Katalogspeichers

Verwenden Sie die Widgets, um einen bestehenden Unity-Catalog-Katalog, ein Schema und ein Volumen zu spezifizieren. Das Notebook speichert den MASSIVE-Cache und die Parquet-Vorhersagen auf diesem Volume. Du brauchst diese Privilegien:

  • USE CATALOG im Katalog und USE SCHEMA im Schema.
  • READ VOLUME und WRITE VOLUME für die Lautstärke.

Jeder MLflow-Lauf schreibt Vorhersagen in sein eigenes Unterverzeichnis unter der konfigurierten Parquet-Ausgabewurzel.

widget_defaults = {
    "uc_catalog": "main",
    "uc_schema": "default",
    "uc_volume": "ray_data",
}
for widget_name, default_value in widget_defaults.items():
    dbutils.widgets.text(widget_name, default_value)

CATALOG = dbutils.widgets.get("uc_catalog")
SCHEMA = dbutils.widgets.get("uc_schema")
VOLUME = dbutils.widgets.get("uc_volume")

volume_path = f"/Volumes/{CATALOG}/{SCHEMA}/{VOLUME}"
parquet_output_root = f"{volume_path}/sgc-raydata-vllm-batch-inference"
massive_cache_path = f"{volume_path}/hf-cache/amazon-massive-1.1"
print(f"Parquet output root: {parquet_output_root}")
print(f"Dataset cache: {massive_cache_path}")

Starten von Ray

ray_init() startet Ray auf der angehängten Compute und druckt die Dashboard-URL für dieses Notizbuch. Die Ray-Verbindung bleibt aktiv, während das Notizbuch verbunden bleibt. Der Actor-Pool verwendet die von Ray gemeldete GPU-Anzahl, sodass jede verfügbare GPU ein vLLM-Modell-Replikat ausführt.

import ray
from serverless_gpu import ray_init

ray_context = ray_init()
ACTOR_COUNT = int(ray.cluster_resources().get("GPU", 0))
if ACTOR_COUNT < 1:
    raise RuntimeError("Ray did not detect a GPU. Attach GPU compute and run the notebook again.")
print(f"Ray detected {ACTOR_COUNT} GPUs; using {ACTOR_COUNT} predictor actors.")

Laden und abtasten von MASSIV

Lade MASSIVE 1.1 in den konfigurierten Cache herunter und wähle dann bei jedem Durchlauf dieselben 2.000 Trainingsbeispiele aus jedem Standort aus. Der erste Standort liefert außerdem die Szenarien- und Absichtsnamen, mit denen die Klassifizierungsaufforderung erstellt wurde.

disable_progress_bar()
download_config = DownloadConfig(cache_dir=f"{massive_cache_path}/downloads")
download_manager = DownloadManager(download_config=download_config)
massive_archive_dir = Path(download_manager.download_and_extract(MASSIVE_ARCHIVE_URL))
massive_data_dir = massive_archive_dir / "1.1" / "data"
locale_datasets = []
scenario_names = None
scenario_intents = None

for locale in LOCALES:
    locale_dataset = load_dataset(
        "json",
        data_files=str(massive_data_dir / f"{locale}.jsonl"),
        split="train",
        cache_dir=f"{massive_cache_path}/datasets",
    )
    locale_dataset = locale_dataset.filter(lambda row: row["partition"] == "train")
    locale_scenarios = sorted(locale_dataset.unique("scenario"))
    if scenario_names is not None and locale_scenarios != scenario_names:
        raise ValueError(f"Scenario labels differ for locale {locale}.")
    if scenario_names is None:
        scenario_names = locale_scenarios
        label_frame = locale_dataset.select_columns(["scenario", "intent"]).to_pandas()
        scenario_intents = {
            scenario: sorted(group["intent"].unique())
            for scenario, group in label_frame.groupby("scenario")
        }
    sample = locale_dataset.shuffle(seed=SEED).select(range(ROWS_PER_LOCALE))
    locale_datasets.append(sample.select_columns(["id", "locale", "utt", "scenario"]))

Erstellen Sie den Ray-Datensatz

Kombinieren Sie die Lokalbeispiele, behalten Sie die für Inferenz und Auswertung benötigten Felder auf und partitionieren Sie die Daten neu, sodass Ray alle Prädiktorakteure beschäftigen kann.

massive_sample = concatenate_datasets(locale_datasets)
records = [
    {
        "input_id": f"{row['locale']}:{row['id']}",
        "locale": row["locale"],
        "utterance": row["utt"],
        "expected_scenario": row["scenario"],
    }
    for row in massive_sample
]
input_dataset = ray.data.from_items(records).repartition(ACTOR_COUNT * 8)
print(f"Prepared {len(records):,} records across {len(LOCALES)} locales and {len(scenario_names)} scenarios.")

Definiere den vLLM-Prädiktor

MASSIVE gruppiert Äußerungen in 18 Szenarien, zum Beispiel alarm, weather und music. Dieses Notizbuch erstellt die erlaubten Labels und die Szenario-zu-Absicht-Anleitung aus dem Datensatz, anstatt sie fest zu kodieren.

Das Szenario-zu-Intention-Mapping hilft Qwen, Labels mit ähnlicher Bedeutung zu unterscheiden. vLLM gibt eines der erlaubten Labels zurück, und ein letzter Normalisierungsschritt markiert jede andere Antwort als ungültig.

scenario_set = set(scenario_names)
scenario_guidance = "\n".join(
    f"- {scenario}: {', '.join(scenario_intents[scenario])}"
    for scenario in scenario_names
)
system_prompt = (
    "Classify the user utterance into exactly one MASSIVE scenario. "
    "Use these scenario-to-intent mappings to distinguish similar labels:\n"
    f"{scenario_guidance}\n"
    "Return only the scenario label."
)

def format_prompt(tokenizer, utterance: str) -> str:
    messages = [
        {"role": "system", "content": system_prompt},
        {"role": "user", "content": utterance},
    ]
    return tokenizer.apply_chat_template(messages, tokenize=False, add_generation_prompt=True)

def normalize_label(response: str) -> str | None:
    normalized = re.sub(r"[^a-z]+", " ", response.lower()).strip()
    return normalized if normalized in scenario_set else None
class VLLMPredictor:
    def __init__(self):
        gpu_ids = ray.get_runtime_context().get_accelerator_ids().get("GPU", [])
        if len(gpu_ids) != 1:
            raise RuntimeError(f"Expected one GPU per actor, but received {gpu_ids}.")
        self.gpu_assignment = str(gpu_ids[0])
        self.llm = LLM(
            model=MODEL_NAME,
            tensor_parallel_size=1,
            dtype="bfloat16",
            max_model_len=MAX_MODEL_LEN,
            max_num_seqs=BATCH_SIZE,
            gpu_memory_utilization=0.90,
            enable_prefix_caching=True,
        )
        self.tokenizer = self.llm.get_tokenizer()
        self.sampling_params = SamplingParams(
            temperature=0.0,
            max_tokens=MAX_OUTPUT_TOKENS,
            structured_outputs=StructuredOutputsParams(choice=scenario_names),
        )

    def __call__(self, batch: pd.DataFrame) -> pd.DataFrame:
        prompts = [format_prompt(self.tokenizer, utterance) for utterance in batch["utterance"]]
        outputs = self.llm.generate(prompts, self.sampling_params, use_tqdm=False)
        raw_responses = [output.outputs[0].text.strip() for output in outputs]
        predicted_scenarios = [normalize_label(response) for response in raw_responses]

        result = batch.copy()
        result["raw_response"] = raw_responses
        # Preserve invalid responses as nulls with a stable string type across batches.
        result["predicted_scenario"] = pd.array(predicted_scenarios, dtype="string")
        result["valid_prediction"] = result["predicted_scenario"].notna()
        result["correct"] = (result["predicted_scenario"] == result["expected_scenario"]).fillna(False)
        result["model_name"] = MODEL_NAME
        result["ray_gpu_assignment"] = self.gpu_assignment
        return result

Batch-Inferenz ausführen und überwachen

VLLMPredictor lädt Qwen einmal, wenn jeder Actor startet, und verwendet dieses Modell dann für jeden erhaltenen Batch wieder. Ray Data startet für jede erkannte GPU einen Actor und weist jeden Batch dem nächsten verfügbaren Actor zu.

Während der Rückschluss ausgeführt wird, öffnen Sie die Ray-Dashboard-URL, die von ray_init() in Zelle 10 gedruckt wird. Nutzen Sie das Dashboard, um die acht Prädiktor-Akteure, GPU-Reservierungen, Aufgabenfortschritte, Logs und Nachzügler zu überprüfen.

predictions = input_dataset.map_batches(
    VLLMPredictor,
    batch_format="pandas",
    batch_size=BATCH_SIZE,
    compute=ray.data.ActorPoolStrategy(size=ACTOR_COUNT),
    num_gpus=1,
)

Materialisieren und die Ergebnisse verfolgen

Ray Data baut diese Pipeline verzögert auf, sodass write_parquet() den Rückschluss ausführt und die Ergebnisse in einem Schritt speichert. Spark liest dann die Parquet-Dateien zur Bewertung, ohne das Modell erneut auszuführen. Die umgebende MLflow-Ausführung erfasst Workload-Parameter, Qualitätsmetriken, Zeitangaben, Durchsatz und Systemmetriken, und Databricks fügt nach Abschluss der Ausführung unter der Zelle einen klickbaren (1 MLflow run)-Link hinzu.

mlflow.set_system_metrics_sampling_interval(2)
with mlflow.start_run(run_name="raydata-massive-qwen25-32b", log_system_metrics=True) as active_run:
    parquet_output_path = f"{parquet_output_root}/{active_run.info.run_id}"
    print(f"Parquet output: {parquet_output_path}")
    mlflow.log_params(
        {
            "model": MODEL_NAME,
            "dataset": DATASET_NAME,
            "dataset_version": "1.1",
            "locales": json.dumps(LOCALES),
            "record_count": len(records),
            "actor_count": ACTOR_COUNT,
            "batch_size": BATCH_SIZE,
            "max_model_len": MAX_MODEL_LEN,
            "max_output_tokens": MAX_OUTPUT_TOKENS,
            "temperature": 0.0,
            "output_constraint": "scenario_choices",
            "system_metrics_interval_seconds": 2,
            "gpu_memory_utilization": 0.90,
        }
    )
    mlflow.set_tags(
        {
            "dataset_source": MASSIVE_ARCHIVE_URL,
            "parquet_output_path": parquet_output_path,
        }
    )

    start_time = time.perf_counter()
    predictions.write_parquet(parquet_output_path)
    cold_start_inclusive_duration_seconds = time.perf_counter() - start_time

    results_df = spark.read.parquet(parquet_output_path)
    aggregate = results_df.agg(
        F.count("*").alias("record_count"),
        F.avg(F.col("correct").cast("double")).alias("overall_accuracy"),
        F.avg(F.col("valid_prediction").cast("double")).alias("valid_prediction_rate"),
        F.countDistinct("ray_gpu_assignment").alias("unique_gpu_assignments"),
    ).first()
    scenario_accuracy_df = results_df.groupBy("expected_scenario").agg(
        F.count("*").alias("record_count"),
        F.avg(F.col("correct").cast("double")).alias("accuracy"),
        F.avg(F.col("valid_prediction").cast("double")).alias("valid_prediction_rate"),
    ).orderBy("expected_scenario")
    macro_scenario_accuracy = scenario_accuracy_df.agg(F.avg("accuracy")).first()[0]
    cold_start_inclusive_records_per_second = (
        aggregate["record_count"] / cold_start_inclusive_duration_seconds
    )
    mlflow.log_metrics(
        {
            "overall_accuracy": aggregate["overall_accuracy"],
            "macro_scenario_accuracy": macro_scenario_accuracy,
            "valid_prediction_rate": aggregate["valid_prediction_rate"],
            "cold_start_inclusive_duration_seconds": cold_start_inclusive_duration_seconds,
            "cold_start_inclusive_records_per_second": cold_start_inclusive_records_per_second,
        }
    )
    mlflow_run_id = active_run.info.run_id

print(f"MLflow run ID: {mlflow_run_id}")
print("Open the '(1 MLflow run)' link attached to this cell for parameters and metrics.")

Das Ergebnis validieren

Die untenstehenden Prüfungen bestätigen, dass die Ausgabe pro Eingabe eine Zeile enthält und dass jeder Prädiktor-Akteur mindestens einen Batch verarbeitet hat.

Das Timing beginnt, bevor Ray die Akteure erstellt und das Modell lädt, daher beinhalten die gemeldete Dauer und Durchsatz die Kaltstartzeit.

if aggregate["record_count"] != len(records):
    raise RuntimeError("The persisted result count does not match the input count.")
if aggregate["unique_gpu_assignments"] != ACTOR_COUNT:
    raise RuntimeError(f"Expected results from {ACTOR_COUNT} Ray GPU assignments.")

print(f"Records: {aggregate['record_count']:,}")
print(f"Overall accuracy: {aggregate['overall_accuracy']:.2%}")
print(f"Macro scenario accuracy: {macro_scenario_accuracy:.2%}")
print(f"Valid prediction rate: {aggregate['valid_prediction_rate']:.2%}")
print(f"Inference duration including actor and model cold start: {cold_start_inclusive_duration_seconds:.1f} seconds")
print(f"Throughput including actor and model cold start: {cold_start_inclusive_records_per_second:.1f} records/second")
print(f"Unique GPU assignments: {aggregate['unique_gpu_assignments']}")

Analysieren Sie die Prädiktionsqualität

Zeigen Sie die Genauigkeit nach Standort, eine Stichprobe von Vorhersagen und die Verteilung der Datensätze über GPU-Aktoren hinweg.

locale_accuracy_df = (
    results_df.groupBy("locale")
    .agg(
        F.count("*").alias("record_count"),
        F.avg(F.col("correct").cast("double")).alias("accuracy"),
        F.avg(F.col("valid_prediction").cast("double")).alias("valid_prediction_rate"),
    )
    .orderBy("locale")
)

print("Accuracy by locale:")
locale_accuracy_df.show(truncate=False)
prediction_columns = [
    "locale", "utterance", "expected_scenario", "predicted_scenario",
    "correct", "ray_gpu_assignment",
]
sample_predictions_df = (
    results_df.select(prediction_columns)
    .orderBy(F.rand(SEED))
    .limit(16)
)
actor_distribution_df = (
    results_df.groupBy("ray_gpu_assignment")
    .agg(F.count("*").alias("record_count"))
    .orderBy("ray_gpu_assignment")
)

print("Sample predictions:")
sample_predictions_df.show(truncate=80)
print("Records by Ray GPU assignment:")
actor_distribution_df.show(truncate=False)

Beispiel-Notebook

Qwen2.5-32B Batch-Inferenz mit Ray Data und vLLM

Notebook abrufen