Classificatietaken met SynapseML

In dit artikel wordt beschreven hoe u een tekstclassificatietaak met twee methoden uitvoert. De ene methode maakt gebruik van gewone pysparkgegevens en de andere methode maakt gebruik van de synapseml bibliotheek. Beide methoden leveren dezelfde prestaties op, maar laten zien hoe SynapseML de codecomplexiteit vermindert ten opzichte van pyspark.

De taak voorspelt of een klantbeoordeling van een boek dat is verkocht op Amazon goed is (beoordeling > 3) of slecht, op basis van de beoordelingstekst. U traint LogisticRegression-cursisten met verschillende hyperparameters en kiest vervolgens het beste model.

Vereiste voorwaarden

  • Maak een notitieblok.
  • Koppel uw notitieblok aan een lakehouse. Selecteer Toevoegen in het linkerdeelvenster in het notitieblok om een bestaand lakehouse toe te voegen of een nieuwe te maken.

Note

Alle bibliotheken die in dit artikel worden gebruikt (pyspark, synapseml, numpy) zijn vooraf geïnstalleerd in de Fabric Spark-runtime. U hoeft geen pakketten te installeren.

De gegevens laden en verkennen

In Fabric notebooks is een Spark-sessie al beschikbaar als de variabele spark. Laad de gegevensset met beoordelingen van Amazon-boeken vanaf een openbare Azure Blob Storage locatie:

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

Controleer of de gegevensset correct is geladen:

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

Functies extraheren en gegevens verwerken

Echte gegevens hebben vaak kenmerken van meerdere typen, zoals tekst, numeriek en categorisch. Als u wilt laten zien hoe u met gemengde functietypen werkt, voegt u twee numerieke functies toe aan de gegevensset: het aantal woorden van de beoordeling en de gemiddelde woordlengte.

Door de gebruiker gedefinieerde functies (UDF's) definiëren

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())

UDF's toepassen met SynapseML UDFTransformer

Gebruik synapseML UDFTransformer om de UDF's in te pakken in pijplijn-compatibele transformatoren:

from synapse.ml.stages import UDFTransformer

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

De functiepijplijn uitvoeren

Pas zowel transformatoren toe als maak een kolom met binaire labels op basis van de classificatie:

from pyspark.ml import Pipeline

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

Controleer de functie-extractie:

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

Classificeren met behulp van pyspark

Als u de beste LogisticRegression-classificatie wilt kiezen met behulp van de pyspark bibliotheek, moet u deze stappen expliciet uitvoeren:

  1. De functies verwerken:
    • De tekstkolom tokeniseren.
    • Hash de getokeniseerde kolom om in een vector met hashing.
    • Voeg de numerieke kenmerken samen met de vector.
  2. Converteer de labelkolom van het type boolean naar het type integer.
  3. Train meerdere LogisticRegression-algoritmen op de train gegevensset met verschillende hyperparameters.
  4. Bereken het gebied onder de ROC-curve (AUC) voor elk getraind model en selecteer het model met de hoogste metrische waarde voor de test gegevensset.
  5. Evalueer het beste model op de validation set.

De gegevens parametriseren en voorbereiden

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())
)

Controleer de gemetriseerde gegevens:

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

Modellen trainen en evalueren

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

Controleer de resultaten:

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

De exacte AUC-waarden zijn afhankelijk van de willekeurige splitsing. Waarden tussen 0,65 en 0,85 verwachten.

Classificeren met SynapseML

De synapseml methode bereikt hetzelfde resultaat met minder stappen. SynapseML verwerkt featurization intern, wat de code vermindert die u moet schrijven:

  1. De TrainClassifier estimator bevat intern de gegevens, zolang de kolommen in de train, testen validation gegevenssets de functies vertegenwoordigen.
  2. De FindBestModel estimator vindt het beste model uit een pool met getrainde modellen door de prestaties van de test gegevensset te evalueren met de opgegeven metrische waarde.
  3. De ComputeModelStatistics transformator berekent meerdere metrische gegevens voor een gescoorde gegevensset (in dit geval de validation gegevensset) tegelijk.
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)
)

Controleer de SynapseML-resultaten:

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

De pyspark- en SynapseML-benaderingen moeten vergelijkbare AUC-waarden produceren, omdat ze hetzelfde modeltype trainen met dezelfde hyperparameters op dezelfde gegevens.

De twee benaderingen vergelijken

Aspect pyspark SynapseML
Verwerking van functies Handmatig (Tokenizer naar HashingTF naar VectorAssembler) Automatisch (verwerkt door TrainClassifier)
Modelselectie Handmatige lus met beoordelaar Ingebouwde FindBestModel
Berekening van metrische gegevens Eén metrische waarde per evaluatiegesprek Meerdere metrische gegevens met ComputeModelStatistics
Regels van code Ongeveer 30 regels Ongeveer 15 regels
Result Dezelfde AUC Dezelfde AUC

Troubleshooting

Probleem Oorzaak Resolutie
AnalysisException: Path does not exist De URL van de openbare blobopslag is tijdelijk niet beschikbaar Wacht enkele minuten en probeer het opnieuw. Connectiviteit controleren door uit te voeren spark.read.parquet("wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet").count()
IllegalArgumentException: Field "features" does not exist De kolomnamen van functies komen niet overeen tussen transformatoren Kolomnamen controleren door uit te voeren data.columns vóór de stap VectorAssembler
NameError: name 'LogisticRegression' is not defined Ontbrekende importstatement Voeg from pyspark.ml.classification import LogisticRegression bovenaan in de cel toe
ModuleNotFoundError: No module named 'synapse.ml' Notebook maakt geen gebruik van Fabric Spark-runtime Controleer of het notebook gebruikmaakt van Fabric Runtime 1.2 of hoger. Selecteer Omgeving op het lint om te controleren.
Lage AUC (lager dan 0,6) Problemen met gegevenssplitsing of convergentie Controleer de labeldistributie met data.groupBy("label").count().show(). U kunt een gegevensset met ongeveer gelijke balans verwachten.
Py4JJavaError: An error occurred while calling interne fout Java/Spark Controleer de Spark-gebruikersinterface op gedetailleerde foutenlogboeken. Start de Spark-sessie opnieuw doorsessiestopsessie te selecteren > en voer vervolgens alle cellen opnieuw uit.

De hulpbronnen opschonen

Als u een nieuw lakehouse voor dit artikel hebt gemaakt en dit niet meer nodig hebt:

  1. Klik in uw werkruimte met de rechtermuisknop op de naam van het lakehouse.
  2. Selecteer Verwijderen.
  3. Bevestig de verwijdering.

Het notitieblok blijft in uw werkruimte, tenzij u het afzonderlijk verwijdert.