Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Usa Qwen2.5-32B-Instruct per classificare 16.000 enunciazioni di assistente vocale multilingue su un runtime AI 8xH100 collegato. Questo quaderno mostra come:
- Costruisci un dataset Ray bilanciato multilingue partendo da MASSIVE 1.1.
- Esegui una replica persistente del modello vLLM su ogni GPU disponibile.
- Monitorare il carico di lavoro con il dashboard Ray e le metriche del sistema MLflow.
- Salva i risultati completi di previsione in formato Parquet in un volume di Unity Catalog.
Note
Questo esempio richiede l'ambiente di IA Databricks versione 5 o superiore.
Connessione alla computazione GPU senza server
- Dal selezionatore di calcolo del notebook, seleziona GPU serverless.
- Nel pannello Ambiente , seleziona l'acceleratore 8xH100 e l'ambiente AI v5 .
- Clicca su Applica, poi conferma l'ambiente.
Il modello Qwen è pubblico e non richiede l'autenticazione di Hugging Face. Il notebook scarica MASSIVE 1.1 dall'archivio pubblico di Amazon.
Importare librerie
AI v5 include i pacchetti Ray, vLLM, Hugging Face Datasets, Transformers, PyTorch e MLflow utilizzati in questo notebook, quindi non è necessaria l'installazione di un pacchetto.
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
Configura il carico di lavoro
Imposta il modello, le località, la dimensione del campione e i parametri di inferenza.
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
Configura l'archiviazione di Unity Catalog
Usa i widget per specificare un catalogo Unity esistente, uno schema e un volume. Il notebook memorizza la cache MASSIVE e le previsioni in formato Parquet in questo volume. Hai bisogno di questi privilegi:
-
USE CATALOGsul catalogo eUSE SCHEMAsullo schema. -
READ VOLUMEeWRITE VOLUMEsul volume.
Ogni run MLflow scrive le previsioni nella propria sottodirectory sotto la radice di uscita configurata di 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}")
Inizio di Ray
ray_init() avvia Ray sulle risorse di calcolo collegate e mostra l'URL della dashboard per questo notebook. La connessione Ray rimane attiva mentre il notebook resta collegato. Il pool di attori utilizza il conteggio delle GPU riportato da Ray, quindi ogni GPU disponibile esegue una replica del modello 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.")
Carica e campiona MASSIVE
Scarica MASSIVE 1.1 nella cache configurata, poi seleziona gli stessi 2.000 esempi di addestramento da ogni località a ogni esecuzione. La prima località fornisce anche i nomi degli scenari e delle intenzioni usati per costruire il prompt di classificazione.
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"]))
Crea il dataset Ray
Combinare i campioni di località, mantenere i campi necessari per l'inferenza e la valutazione, e ripartizionare i dati in modo che Ray possa tenere occupati tutti gli attori predittori.
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.")
Definisci il predittore vLLM
MASSIVE raggruppa enunciati in 18 scenari, come alarm, weather, e music. Questo notebook genera le etichette consentite e le linee guida per l'associazione tra scenari e intenti a partire dal dataset, invece di definirle in modo statico.
La mappatura scenario-intenzione aiuta Qwen a distinguere etichette con significati simili. vLLM restituisce una delle etichette consentite e un passaggio finale di normalizzazione segna qualsiasi altra risposta come invalida.
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
Eseguire e monitorare l'inferenza batch
VLLMPredictor carica Qwen una volta quando ogni attore inizia, poi riutilizza quel modello per ogni batch che riceve. Ray Data avvia un attore per ogni GPU rilevata e programma ogni batch sull'attore disponibile successivo.
Mentre è in esecuzione l'inferenza, apri l'URL della dashboard di Ray stampato da ray_init() nella cella 10. Usa la dashboard per ispezionare gli otto attori predittivi, le prenotazioni GPU, i progressi dei compiti, i log e i ritardatari.
predictions = input_dataset.map_batches(
VLLMPredictor,
batch_format="pandas",
batch_size=BATCH_SIZE,
compute=ray.data.ActorPoolStrategy(size=ACTOR_COUNT),
num_gpus=1,
)
Materializza e monitora i risultati
Ray Data costruisce questa pipeline in modo lazy, quindi write_parquet() esegue l'inferenza e salva i risultati in un unico passaggio. Spark poi legge i file Parquet per la valutazione senza eseguire nuovamente il modello. L'esecuzione MLflow associata acquisisce i parametri del carico di lavoro, le metriche di qualità, i tempi di esecuzione, il throughput e le metriche di sistema, e Databricks aggiunge un (1 MLflow run) link cliccabile sotto la cella quando termina.
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.")
Convalidare i risultati
I controlli seguenti confermano che l'output contiene una riga per ogni input e che ogni attore predittore ha gestito almeno un lotto.
Il timing inizia prima che Ray crei gli attori e carichi il modello, quindi la durata e la produttività riportate includono il tempo di avvio a freddo.
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']}")
Analizza la qualità delle previsioni
Mostra l'accuratezza per località, un campione di previsioni e la distribuzione dei record tra gli attori 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)