Luokitustehtävät SynapseML:n avulla

Tässä artikkelissa näytetään, miten tekstin luokittelutehtävä suoritetaan kahdella menetelmällä. Toinen menetelmä käyttää pelkkää pyspark, ja toinen kirjastoa synapseml . Molemmat menetelmät tuottavat saman suorituskyvyn, mutta korostavat, miten SynapseML vähentää koodin monimutkaisuutta verrattuna .pyspark

Tehtävä ennustaa, onko Amazonissa myydyn kirjan asiakasarvostelu hyvä vai > huono arvostelutekstin perusteella. Koulutat LogisticRegression-oppijoita eri hyperparametreilla ja valitset sitten parhaan mallin.

Edellytykset

  • Luo muistikirja.
  • Liitä muistikirjasi Lakehouseen. Muistikirjassa valitse Lisää vasemmalta ruudulta liittääksesi olemassa olevan järvenrakennuksen tai luodaksesi uuden.

Muistio

Kaikki tässä artikkelissa käytetyt kirjastot (pyspark, synapseml, numpy) on valmiiksi asennettu Fabric Sparkin ajonaikaan. Sinun ei tarvitse asentaa mitään paketteja.

Lataa ja tutki dataa

Fabric muistikirjoissa Spark-istunto on jo saatavilla spark muuttujana. Lataa Amazonin kirja-arvosteluaineisto julkisesta Azure Blob Storage -sijainnista:

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

Varmista, että aineisto latautuu oikein:

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

Poimi ominaisuuksia ja prosessitietoja

Todellisessa datassa on usein useita piirteitä, esimerkiksi teksti, numeerinen ja kategorinen. Havainnollistaaksesi työskentelyä sekoitettujen ominaisuustyyppien kanssa, lisää aineistoon kaksi numeerista ominaisuutta: arvostelun sanamäärä ja keskimääräinen sananpituus.

Määrittele käyttäjän määrittelemät funktiot (UDF:t)

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

Sovella UDF:iä SynapseML:llä UDFTransformer

Käytä from SynapseML: UDFTransformer ää kääriäksesi UDF:t putkistoyhteensopiviin muuntajiin:

from synapse.ml.stages import UDFTransformer

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

Suorita ominaisuusputki

Käytä molempia muuntajia ja luo binäärinen tunnistesarake luokituksesta:

from pyspark.ml import Pipeline

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

Vahvista ominaisuuksien poimiminen:

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

Luokittele pysparkin avulla

Parhaan LogisticRegression-luokittelun pyspark valitsemiseksi kirjaston avulla sinun on suoritettava nämä vaiheet erikseen:

  1. Prosessoi ominaisuuksia:
    • Tokenisoi tekstisarake.
    • Hajauta tokenisoitu sarake vektoriksi hajautuksella.
    • Yhdistä numeeriset ominaisuudet vektoriin.
  2. Heitä label-sarake totuusarvosta kokonaislukutyyppiin.
  3. Kouluta useita LogisticRegression-algoritmeja aineistolla train eri hyperparametreilla.
  4. Laske ROC-käyrän (AUC) alla oleva pinta-ala jokaiselle koulutetulle mallille ja valitse malli, jolla on datan korkein metriikka test .
  5. Arvioi paras malli lavasteissa validation .

Ominaisuudet ja tiedot valmisteltu

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

Varmista ominaisuuksilla varustetut tiedot:

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

Kouluta ja arvioi malleja

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

Varmista tulokset:

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

Muistio

Tarkat AUC-arvot riippuvat satunnaisesta jaosta. Odota arvoja välillä 0,65–0,85.

Luokittele SynapseML:n avulla

Lähestymistapa synapseml saavuttaa saman tuloksen vähemmillä askelilla. SynapseML hoitaa ominaisuuksien käsittelyn sisäisesti, mikä vähentää kirjoittamiseen tarvittavaa koodia:

  1. Estimaattori TrainClassifier esittelee datan sisäisesti, kunhan sarakkeet , traintest, ja validation aineistot edustavat ominaisuuksia.
  2. Estimaattori FindBestModel löytää parhaan mallin koulutettujen mallien joukosta arvioimalla aineiston suorituskykyä test määritellyn mittarin avulla.
  3. Muuntaja ComputeModelStatistics laskee useita mittareita pisteytetystä aineistosta (tässä tapauksessa aineistosta validation ) samanaikaisesti.
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)
)

Varmista SynapseML-tulokset:

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

Muistio

Pyspark- ja SynapseML-lähestymistavat tuottavat samankaltaisia AUC-arvoja, koska ne kouluttavat samaa mallityyppiä samoilla hyperparametreilla samalla datalla.

Vertaa näitä kahta lähestymistapaa

Ominaisuus Pyspark SynapseML
Ominaisuuksien käsittely Manuaali (Tokenizerista HashingTF:hen VectorAssembleriin) Automaattinen (hoitaa TrainClassifier)
Mallin valinta Manuaalinen silmukka arvioijan kanssa Sisäänrakennettu FindBestModel
Metriikan laskenta Yksi metriikka per arviointikutsu Useita mittareita ComputeModelStatistics
Koodirivit Noin 30 linjaa Noin 15 linjaa
Tulos Sama AUC Sama AUC

Vianmääritys

Ongelma Syy Ratkaisu
AnalysisException: Path does not exist Julkisen blobin tallennus-URL on tilapäisesti poissa käytöstä Odota muutama minuutti ja yritä uudelleen. Varmista yhteys ajamalla spark.read.parquet("wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet").count()
IllegalArgumentException: Field "features" does not exist Ominaisuussarakkeiden nimet eivät täsmää muuntajien välillä Varmista sarakkeen nimet ajamalla data.columns ennen VectorAssembler-vaihetta
NameError: name 'LogisticRegression' is not defined Puuttuva tuontilauseke Lisää from pyspark.ml.classification import LogisticRegression solun yläosaan
ModuleNotFoundError: No module named 'synapse.ml' Notebook ei käytä Fabric Spark -ajonaikaa Varmista, että kannettava käyttää Fabric Runtime 1.2:ta tai uudempaa. Valitse nauhasta Ympäristö tarkistaaksesi.
Alhainen AUC (alle 0,6) Datan jakautumisongelma tai konvergenssiongelmat Varmista etikettijakauma .data.groupBy("label").count().show() Odota karkeasti tasapainotettua aineistoa.
Py4JJavaError: An error occurred while calling Java/Sparkin sisäinen virhe Tarkista Sparkin käyttöliittymästä yksityiskohtaiset virhelokit. Käynnistä Spark-istunto uudelleen valitsemalla Session>Stop -istunto, ja suorita sitten kaikki solut uudestaan.

Puhdista resurssit

Jos loit uuden järvimajan tätä artikkelia varten etkä enää tarvitse sitä:

  1. Työtilassasi napsauta järventalon nimeä hiiren oikealla.
  2. Valitse Poista.
  3. Vahvista poistaminen.

Muistikirja pysyy työtilassasi, ellei sitä poisteta erikseen.