Wnioskowanie wsadowe Qwen2.5-32B przy użyciu Ray Data i vLLM

Użyj Qwen2.5-32B-Instruct, aby sklasyfikować 16 000 wielojęzycznych wypowiedzi asystenta głosowego na dołączonym środowisku wykonawczym 8xH100 AI. Ten notatnik pokazuje, jak:

  • Zbuduj zrównoważony, wielojęzyczny zestaw danych promieni z MASSIVE 1.1.
  • Uruchom jedną trwałą replikę modelu vLLM na każdej dostępnej kartie graficznej.
  • Monitoruj obciążenie za pomocą pulpitu Ray i metryk systemu MLflow.
  • Zapisz pełne wyniki predykcji jako plik Parquet w woluminie Unity Catalog.

Note

Ten przykład wymaga środowiska AI Databricks w wersji 5 lub wyższej.

Nawiązywanie połączenia z bezserwerowym przetwarzaniem GPU

  1. Z selektora obliczeń notatnika wybierz GPU bez serwera.
  2. W panelu Środowisko wybierz akcelerator 8xH100 oraz środowisko AI v5 .
  3. Kliknij Zastosuj, a następnie potwierdź środowisko.

Model Qwen jest publiczny i nie wymaga uwierzytelniania Hugging Face. Notebook pobiera MASSIVE 1.1 z publicznego archiwum Amazona.

Importowanie bibliotek

AI v5 zawiera pakiety Ray, vLLM, Hugging Face Datasets, Transformers, PyTorch i MLflow używane w tym notatniku, więc nie trzeba instalować żadnych pakietów.

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

Konfiguruj obciążenie

Ustaw model, lokalizacje, wielkość próby oraz parametry wnioskowania.

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

Konfiguruj pamięć danych katalogowych Unity

Użyj widżetów, aby określić istniejący katalog, schemat i wolumin w Unity Catalog. Notatnik przechowuje dużą pamięć podręczną oraz predykcje w formacie Parquet w tym woluminie. Potrzebujesz następujących przywilejów:

  • USE CATALOG w katalogu i USE SCHEMA w schemacie.
  • READ VOLUME i WRITE VOLUME na regulatorze głośności.

Każde uruchomienie MLflow zapisuje predykcje we własnym podkatalogu w skonfigurowanym katalogu głównym danych wyjściowych Parquet.

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}")

Uruchamianie promienia

ray_init() uruchamia Ray na dołączonych zasobach obliczeniowych i wyświetla adres URL pulpitu nawigacyjnego dla tego notatnika. Połączenie Ray pozostaje aktywne, podczas gdy notatnik pozostaje połączony. Pula aktorów korzysta z liczby GPU raportowanej przez Ray, więc każdy dostępny GPU uruchamia jedną replikę modelu vLLM.

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.")

Załaduj i próbkuj MASSIVE

Pobierz MASSIVE 1.1 do skonfigurowanej pamięci podręcznej, a następnie wybierz te same 2000 przykładów treningowych z każdego miejsca przy każdym uruchomieniu. Pierwsza lokalizacja zawiera także nazwy scenariuszy i intencji użyte do budowy promptu klasyfikacji.

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"]))

Utwórz zbiór danych promieni

Połącz próbki dla poszczególnych ustawień regionalnych, zachowaj pola potrzebne do wnioskowania i oceny oraz ponownie podziel dane, aby Ray mógł w pełni wykorzystać wszystkie aktory wykonujące predykcję.

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.")

Zdefiniuj predyktor vLLM

MASSIVE grupuje wypowiedzi w 18 scenariuszy, takich jak alarm, weather, oraz music. Ten notatnik tworzy dozwolone etykiety oraz wskazówki dotyczące przypisywania scenariuszy do intencji na podstawie zbioru danych, zamiast wpisywać je na stałe.

Mapowanie scenariuszy na intencje pomaga modelowi Qwen odróżniać etykiety o podobnym znaczeniu. vLLM zwraca jedną z dozwolonych etykiet, a ostatni krok normalizacji oznacza każdą inną odpowiedź jako nieważną.

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

Uruchamiaj i monitoruj wnioskowanie wsadowe

VLLMPredictor Qwen ładuje się raz przy każdym aktorze, a następnie ponownie używa tego modelu dla każdej otrzymanej partii. Ray Data uruchamia jednego aktora na wykryty GPU i planuje każdą partię na następnym dostępnym aktorze.

Podczas wykonywania wnioskowania otwórz adres URL dashboardu Ray wydrukowany w ray_init() Cell 10. Użyj panelu, aby sprawdzić ośmiu aktorów predyktora, rezerwacje GPU, postęp zadań, dzienniki i opóźnione zadania.

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

Materializuj i śledź wyniki

Ray Data buduje ten pipeline leniwie, więc write_parquet() uruchamia wnioskowanie i zapisuje wyniki w jednym kroku. Spark następnie odczytuje pliki Parquet do oceny bez ponownego uruchamiania modelu. Otaczający proces MLflow rejestruje parametry obciążenia, metryki jakości, czas, przepustowość i metryki systemowe, a Databricks dodaje klikalny (1 MLflow run) link poniżej komórki po zakończeniu.

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.")

Sprawdzanie poprawności wyników

Poniższe kontrole potwierdzają, że wyjście zawiera jeden wiersz na wejście oraz że każdy aktor predyktora obsłużył co najmniej jedną partię.

Czas rozpoczyna się przed utworzeniem aktorów przez Raya i załadowaniem modelu, więc raportowany czas trwania i przepustowość obejmują czas zimnego startu.

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']}")

Analiza jakości predykcji

Pokaż dokładność według lokalizacji, próbkę przewidywań oraz rozkład rekordów między aktorami GPU.

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)

Przykładowy notatnik

Wnioskowanie wsadowe Qwen2.5-32B przy użyciu Ray Data i vLLM

Pobierz laptopa