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.
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.
- 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:
- Przetwórz funkcje:
- Tokenizowanie kolumny tekstowej.
- Przekształć kolumnę po tokenizacji do wektora za pomocą haszowania.
- Połącz cechy liczbowe z wektorem.
- Rzutuj kolumnę etykiet z typu logicznego na typ liczby całkowitej.
- Wytrenuj wiele algorytmów LogisticRegression na zbiorze danych
trainz różnymi hiperparametrami. - Oblicz obszar pod krzywą ROC (AUC) dla każdego wytrenowanego modelu i wybierz model z najwyższą metryką w
testzestawie danych. - Oceń najlepszy model na
validationzestawie.
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ć:
- Estymator
TrainClassifierautomatycznie wyodrębnia cechy z danych, pod warunkiem że kolumny w zestawach danychtrain,testivalidationreprezentują cechy. - Narzędzie
FindBestModeldo szacowania znajduje najlepszy model z puli wytrenowanych modeli, oceniając wydajność zestawutestdanych za pomocą określonej metryki. - Transformator
ComputeModelStatisticsoblicza jednocześnie wiele metryk dla ocenionego zbioru danych (w tym przypadku zbioru danychvalidation).
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:
- W obszarze roboczym kliknij prawym przyciskiem myszy nazwę lakehouse.
- Wybierz opcję Usuń.
- Potwierdź usunięcie.
Notatnik pozostaje w obszarze roboczym, chyba że usuniesz go osobno.