Zadania klasyfikacji przy użyciu usługi SynapseML

W tym artykule pokazano, jak wykonać zadanie klasyfikacji tekstu przy użyciu dwóch metod. Jedna metoda używa zwykłego pyspark, a druga używa biblioteki synapseml. Obie metody dają taką samą wydajność, ale podkreślają, jak usługa SynapseML zmniejsza złożoność kodu w porównaniu z pyspark.

Zadanie przewiduje, czy recenzja książki sprzedanej przez klienta na Amazon jest dobra (ocena > 3) czy zła na podstawie tekstu przeglądu. Trenujesz uczniów usługi LogisticsRegression przy użyciu różnych hiperparametrów, a następnie wybierasz najlepszy model.

Wymagania wstępne

  • Uzyskaj subskrypcję usługi Microsoft Fabric. Możesz też utworzyć konto bezpłatnej wersji próbnej usługi Microsoft Fabric.

  • Zaloguj się do Microsoft Fabric.

  • Przełącz na Fabric, używając przełącznika doświadczenia w dolnym lewym rogu twojej strony głównej.

    Zrzut ekranu przedstawiający wybór Fabric w menu przełącznika środowiska.

  • Utwórz notatnik.
  • Dołącz swój notebook do lakehouse. W notesie wybierz pozycję Dodaj w lewym panelu, aby dołączyć istniejący lakehouse lub utworzyć nowy.

Note

Wszystkie biblioteki używane w tym artykule (pyspark, synapseml, numpy) są wstępnie zainstalowane w środowisku uruchomieniowym Fabric Spark. Nie trzeba instalować żadnych pakietów.

Ładowanie i eksplorowanie danych

W notesach Fabric sesja Spark jest już dostępna jako zmienna spark. Załaduj zestaw danych przeglądów książek amazon z publicznej lokalizacji Azure Blob Storage:

rawData = spark.read.parquet(
    "wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet"
)
rawData.show(5)

Sprawdź, czy zestaw danych został załadowany poprawnie:

print(f"Row count: {rawData.count()}")
print(f"Columns: {rawData.columns}")
assert rawData.count() == 10000, "Expected 10,000 rows"
assert set(rawData.columns) == {"text", "rating"}, "Expected columns: text, rating"
print("Data loaded successfully")

Wyodrębnianie funkcji i przetwarzanie danych

Rzeczywiste dane często mają funkcje wielu typów, na przykład tekst, liczbowe i podzielone na kategorie. Aby zademonstrować pracę z typami funkcji mieszanych, dodaj do zestawu danych dwie funkcje liczbowe: liczbę wyrazów recenzji i średnią długość słowa.

Definiowanie funkcji zdefiniowanych przez użytkownika (UDF)

from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType, DoubleType
import numpy as np


def calc_word_count(s):
    return len(s.split())


def calc_word_length(s):
    ss = [len(w) for w in s.split()]
    return round(float(np.mean(ss)), 2)


wordLengthUDF = udf(calc_word_length, DoubleType())
wordCountUDF = udf(calc_word_count, IntegerType())

Stosowanie funkcji zdefiniowanych przez użytkownika za pomocą składnika SynapseML UDFTransformer

Użyj elementu UDFTransformer z SynapseML, aby opakować funkcje UDF w transformatory zgodne z potokiem:

from synapse.ml.stages import UDFTransformer

wordLengthTransformer = UDFTransformer(
    inputCol="text", outputCol="wordLength", udf=wordLengthUDF
)
wordCountTransformer = UDFTransformer(
    inputCol="text", outputCol="wordCount", udf=wordCountUDF
)

Uruchamianie potoku funkcji

Zastosuj obie transformacje i utwórz kolumnę etykiety binarnej na podstawie oceny:

from pyspark.ml import Pipeline

data = (
    Pipeline(stages=[wordLengthTransformer, wordCountTransformer])
    .fit(rawData)
    .transform(rawData)
    .withColumn("label", rawData["rating"] > 3)
    .drop("rating")
)

Sprawdź ekstrakcję cech:

data.show(5)
print(f"Columns: {data.columns}")
assert "wordLength" in data.columns, "wordLength column missing"
assert "wordCount" in data.columns, "wordCount column missing"
assert "label" in data.columns, "label column missing"
assert "rating" not in data.columns, "rating column should be dropped"
print("Feature extraction successful")

Klasyfikowanie przy użyciu narzędzia pyspark

Aby wybrać najlepszy klasyfikator LogisticsRegression przy użyciu pyspark biblioteki, należy jawnie wykonać następujące kroki:

  1. Przetwórz funkcje:
    • Tokenizowanie kolumny tekstowej.
    • Przekształć kolumnę po tokenizacji do wektora za pomocą haszowania.
    • Połącz cechy liczbowe z wektorem.
  2. Rzutuj kolumnę etykiet z typu logicznego na typ liczby całkowitej.
  3. Wytrenuj wiele algorytmów LogisticRegression na zbiorze danych train z różnymi hiperparametrami.
  4. Oblicz obszar pod krzywą ROC (AUC) dla każdego wytrenowanego modelu i wybierz model z najwyższą metryką w test zestawie danych.
  5. Oceń najlepszy model na validation zestawie.

Tworzenie cech i przygotowanie danych

from pyspark.ml.feature import Tokenizer, HashingTF, VectorAssembler
from pyspark.sql.types import IntegerType

# Tokenize the text column
tokenizer = Tokenizer(inputCol="text", outputCol="tokenizedText")
numFeatures = 10000
hashingScheme = HashingTF(
    inputCol="tokenizedText", outputCol="TextFeatures", numFeatures=numFeatures
)
tokenizedData = tokenizer.transform(data)
featurizedData = hashingScheme.transform(tokenizedData)

# Merge text and numeric features into one feature column
featureColumnsArray = ["TextFeatures", "wordCount", "wordLength"]
assembler = VectorAssembler(inputCols=featureColumnsArray, outputCol="features")
assembledData = assembler.transform(featurizedData)

# Select only the label and features columns, cast label to integer
processedData = assembledData.select("label", "features").withColumn(
    "label", assembledData.label.cast(IntegerType())
)

Sprawdź dane cechowane:

print(f"Feature vector size: {processedData.first()['features'].size}")
print(f"Label values: {sorted(processedData.select('label').distinct().rdd.flatMap(lambda x: x).collect())}")
assert processedData.first()["features"].size == 10002, "Expected 10000 text + 2 numeric features"
print("Featurization successful")

Trenowanie i ocenianie modeli

from pyspark.ml.evaluation import BinaryClassificationEvaluator
from pyspark.ml.classification import LogisticRegression

# Split the data into train, test, and validation sets
train, test, validation = processedData.randomSplit([0.60, 0.20, 0.20], seed=123)

# Train models with different regularization parameters
lrHyperParams = [0.05, 0.1, 0.2, 0.4]
logisticRegressions = [
    LogisticRegression(regParam=hyperParam) for hyperParam in lrHyperParams
]
evaluator = BinaryClassificationEvaluator(
    rawPredictionCol="rawPrediction", metricName="areaUnderROC"
)
metrics = []
models = []

# Train each model and evaluate on the test set
for learner in logisticRegressions:
    model = learner.fit(train)
    models.append(model)
    scoredData = model.transform(test)
    metrics.append(evaluator.evaluate(scoredData))

bestMetric = max(metrics)
bestModel = models[metrics.index(bestMetric)]

# Evaluate the best model on the validation dataset
scoredVal = bestModel.transform(validation)
validationAUC = evaluator.evaluate(scoredVal)
print(f"Best model's AUC on validation set = {validationAUC:.4f}")

Sprawdź wyniki:

print(f"Number of models trained: {len(models)}")
print(f"Best regularization parameter: {lrHyperParams[metrics.index(bestMetric)]}")
print(f"Test AUC scores: {[f'{m:.4f}' for m in metrics]}")
assert 0.5 < validationAUC <= 1.0, f"AUC {validationAUC} is outside expected range (0.5, 1.0]"
print(f"pyspark classification complete - AUC: {validationAUC:.4f}")

Note

Dokładne wartości AUC zależą od losowego podziału. Oczekiwano wartości z zakresu od 0,65 do 0,85.

Klasyfikowanie przy użyciu usługi SynapseML

Podejście synapseml osiąga ten sam wynik z mniejszą liczbą kroków. Usługa SynapseML obsługuje cechowanie wewnętrznie, co zmniejsza kod, który należy napisać:

  1. Estymator TrainClassifier automatycznie wyodrębnia cechy z danych, pod warunkiem że kolumny w zestawach danych train, test i validation reprezentują cechy.
  2. Narzędzie FindBestModel do szacowania znajduje najlepszy model z puli wytrenowanych modeli, oceniając wydajność zestawu test danych za pomocą określonej metryki.
  3. Transformator ComputeModelStatistics oblicza jednocześnie wiele metryk dla ocenionego zbioru danych (w tym przypadku zbioru danych validation).
from synapse.ml.train import TrainClassifier, ComputeModelStatistics
from synapse.ml.automl import FindBestModel
from pyspark.ml.classification import LogisticRegression

# Split the raw feature data (SynapseML handles featurization internally)
train, test, validation = data.randomSplit([0.60, 0.20, 0.20], seed=123)

# Train models with different regularization parameters
lrHyperParams = [0.05, 0.1, 0.2, 0.4]
logisticRegressions = [
    LogisticRegression(regParam=hyperParam) for hyperParam in lrHyperParams
]
lrmodels = [
    TrainClassifier(model=lrm, labelCol="label", numFeatures=10000).fit(train)
    for lrm in logisticRegressions
]

# Select the best model based on AUC
bestModel = FindBestModel(evaluationMetric="AUC", models=lrmodels).fit(test)

# Compute metrics on the validation dataset
predictions = bestModel.transform(validation)
metrics = ComputeModelStatistics().transform(predictions)
print(
    "Best model's AUC on validation set = "
    + "{0:.2f}%".format(metrics.first()["AUC"] * 100)
)

Sprawdź wyniki usługi SynapseML:

auc_value = metrics.first()["AUC"]
print(f"Available metrics: {metrics.columns}")
assert 0.5 < auc_value <= 1.0, f"AUC {auc_value} is outside expected range (0.5, 1.0]"
print(f"SynapseML classification complete - AUC: {auc_value:.4f}")

Note

Metody pyspark i SynapseML powinny generować podobne wartości AUC, ponieważ trenują ten sam typ modelu przy użyciu tych samych hiperparametrów na tych samych danych.

Porównanie dwóch podejść

Aspekt pyspark SynapseML
Przetwarzanie cech Ręcznie (Tokenizer do HashingTF do VectorAssembler) Automatyczne (obsługiwane przez TrainClassifier)
Wybór modelu Ręczna pętla z modułem oceny Wbudowane FindBestModel
Obliczenia metryk Jedna metryka na każde wywołanie oceny Wiele wskaźników z ComputeModelStatistics
Wiersze kodu Około 30 wierszy Około 15 wierszy
Result To samo AUC To samo AUC

Troubleshooting

Problematyka Przyczyna Resolution
AnalysisException: Path does not exist Publiczny adres URL magazynu obiektów blob jest tymczasowo niedostępny Odczekaj kilka minut i spróbuj ponownie. Zweryfikuj łączność, uruchamiając polecenie spark.read.parquet("wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet").count()
IllegalArgumentException: Field "features" does not exist Nazwy kolumn funkcji nie są zgodne między transformatorami Zweryfikuj nazwy kolumn, uruchamiając data.columns polecenie przed krokiem VectorAssembler
NameError: name 'LogisticRegression' is not defined Brak instrukcji importu Dodaj from pyspark.ml.classification import LogisticRegression w górnej części komórki
ModuleNotFoundError: No module named 'synapse.ml' Notes nie używa środowiska uruchomieniowego platformy Spark Fabric Sprawdź, czy notes korzysta ze środowiska uruchomieniowego Fabric Runtime 1.2 lub nowszego. Wybierz pozycję Środowisko na wstążce, aby sprawdzić.
Niska AUC (poniżej 0,6) Problem z podziałem danych lub problemy zbieżności Sprawdź dystrybucję etykiet za pomocą polecenia data.groupBy("label").count().show(). Oczekiwano mniej więcej zrównoważonego zestawu danych.
Py4JJavaError: An error occurred while calling błąd wewnętrzny Java/Spark Sprawdź interfejs użytkownika platformy Spark, aby uzyskać szczegółowe dzienniki błędów. Uruchom ponownie sesję Spark, wybierając Sesja>Zatrzymaj sesję, a następnie ponownie uruchom wszystkie komórki.

Uprzątnij zasoby

Jeśli utworzysz nowy element lakehouse na potrzeby tego artykułu i już go nie potrzebujesz:

  1. W obszarze roboczym kliknij prawym przyciskiem myszy nazwę lakehouse.
  2. Wybierz opcję Usuń.
  3. Potwierdź usunięcie.

Notatnik pozostaje w obszarze roboczym, chyba że usuniesz go osobno.