Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
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
- Z selektora obliczeń notatnika wybierz GPU bez serwera.
- W panelu Środowisko wybierz akcelerator 8xH100 oraz środowisko AI v5 .
- 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 CATALOGw katalogu iUSE SCHEMAw schemacie. -
READ VOLUMEiWRITE VOLUMEna 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)