Tâches de classification à l’aide de SynapseML

Cet article montre comment effectuer une tâche de classification de texte avec deux méthodes. Une méthode utilise une méthode simple pysparket l’autre utilise la synapseml bibliothèque. Les deux méthodes offrent les mêmes performances, mais mettent en évidence la façon dont SynapseML réduit la complexité du code par rapport à pyspark.

La tâche prédit si un avis client d’un livre vendu sur Amazon est bon (évaluation > 3) ou mauvais, en fonction du texte de révision. Vous entraînez les apprenants LogisticRegression avec différents hyperparamètres, puis choisissez le meilleur modèle.

Conditions préalables

  • Créez un bloc-notes.
  • Fixez votre notebook à un lakehouse. Dans le bloc-notes, sélectionnez Ajouter dans le volet gauche pour attacher un lakehouse existant ou en créer un nouveau.

Note

Toutes les bibliothèques utilisées dans cet article (pyspark, synapseml, numpy) sont préinstallées dans le runtime Fabric Spark. Vous n’avez pas besoin d’installer de packages.

Charger et explorer les données

Dans Fabric notebooks, une session Spark est déjà disponible en tant que variable spark. Chargez le jeu de données des avis sur les livres Amazon à partir d’un emplacement public dans Stockage Blob Azure :

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

Vérifiez que le jeu de données est correctement chargé :

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

Extraire des fonctionnalités et traiter des données

Les données réelles ont souvent des caractéristiques de plusieurs types, par exemple, du texte, des données numériques et des catégories. Pour illustrer l’utilisation de types de caractéristiques mixtes, ajoutez deux caractéristiques numériques au jeu de données : le nombre de mots de la révision et la longueur moyenne du mot.

Définir des fonctions définies par l’utilisateur (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())

Appliquer des fonctions définies par l’utilisateur avec SynapseML UDFTransformer

Utilisez le UDFTransformer de SynapseML pour intégrer les UDF dans des transformateurs compatibles avec les pipelines :

from synapse.ml.stages import UDFTransformer

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

Exécuter le pipeline de fonctionnalités

Appliquez les deux transformateurs et créez une colonne d’étiquette binaire à partir de l’évaluation :

from pyspark.ml import Pipeline

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

Vérifiez l’extraction de fonctionnalités :

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

Classifier à l’aide de pyspark

Pour choisir le classifieur LogisticRegression le mieux adapté à l’aide de la pyspark bibliothèque, vous devez effectuer explicitement les étapes suivantes :

  1. Traiter les fonctionnalités :
    • Tokenisez la colonne de texte.
    • Hachez la colonne tokenisée en vecteur par hachage.
    • Fusionnez les caractéristiques numériques avec le vecteur.
  2. Convertissez la colonne d’étiquette du type booléen en type entier.
  3. Entraîner plusieurs algorithmes LogisticRegression sur le train jeu de données avec différents hyperparamètres.
  4. Calculez la zone sous la courbe ROC (AUC) pour chaque modèle entraîné et sélectionnez le modèle avec la métrique la plus élevée sur le test jeu de données.
  5. Évaluez le meilleur modèle sur l’ensemble validation .

Caractérisation et préparation des données

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

Vérifiez les données caractérisations :

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

Entraîner et évaluer des modèles

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

Vérifiez les résultats :

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

Les valeurs AUC exactes dépendent du fractionnement aléatoire. Attendez-vous aux valeurs comprises entre 0,65 et 0,85.

Classifier à l’aide de SynapseML

L’approche synapseml obtient le même résultat avec moins d’étapes. SynapseML gère la caractérisation en interne, ce qui réduit le code que vous devez écrire :

  1. L’estimateur TrainClassifier transforme les données en caractéristiques en interne, à condition que les colonnes des jeux de données train, test et validation correspondent aux caractéristiques.
  2. L’estimateur FindBestModel trouve le meilleur modèle à partir d’un pool de modèles formés en évaluant les performances sur le test jeu de données avec la métrique spécifiée.
  3. Le ComputeModelStatistics transformateur calcule plusieurs métriques sur un jeu de données noté (dans ce cas, le validation jeu de données) en même temps.
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)
)

Vérifiez les résultats 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

Les approches pyspark et SynapseML doivent produire des valeurs AUC similaires, car elles entraînent le même type de modèle avec les mêmes hyperparamètres sur les mêmes données.

Comparer les deux approches

Aspect pyspark SynapseML
Traitement des fonctionnalités Manuel (Tokenizer vers HashingTF vers VectorAssembler) Automatique (géré par TrainClassifier)
Sélection du modèle Boucle manuelle avec évaluateur Intégré FindBestModel
Calcul des métriques Métrique unique par appel d’évaluation Plusieurs métriques avec ComputeModelStatistics
Lignes de code Environ 30 lignes Environ 15 lignes
Résultat Même AUC Même AUC

Résolution des problèmes

Problème Cause Résolution
AnalysisException: Path does not exist L’URL publique du stockage Blob est temporairement indisponible Veuillez patienter quelques minutes et réessayez. Vérifier la connectivité en exécutant spark.read.parquet("wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet").count()
IllegalArgumentException: Field "features" does not exist Les noms de colonnes de fonctionnalité ne correspondent pas entre les transformateurs Vérifier les noms des colonnes en exécutant data.columns avant l’étape VectorAssembler
NameError: name 'LogisticRegression' is not defined Instruction d’import manquante Ajouter from pyspark.ml.classification import LogisticRegression en haut de la cellule
ModuleNotFoundError: No module named 'synapse.ml' Notebook n'utilise pas le runtime Spark de Fabric Vérifiez que le notebook utilise Fabric Runtime 1.2 ou version ultérieure. Sélectionnez Environnement dans le ruban à vérifier.
Faible AUC (inférieur à 0,6) Problèmes de fractionnement des données ou de convergence Vérifiez la distribution des étiquettes avec data.groupBy("label").count().show(). Attendez-vous à un jeu de données à peu près équilibré.
Py4JJavaError: An error occurred while calling erreur interne Java/Spark Vérifiez l’interface utilisateur Spark pour obtenir des journaux d’erreurs détaillés. Redémarrez la session Spark en sélectionnant Session>, puis réexécutez toutes les cellules.

Nettoyer les ressources

Si vous avez créé un nouveau lakehouse pour cet article et que vous n’en avez plus besoin :

  1. Dans votre espace de travail, faites un clic droit sur le nom du lakehouse.
  2. Sélectionnez Supprimer.
  3. Confirmez la suppression.

Le bloc-notes reste dans votre espace de travail, sauf si vous le supprimez séparément.