Koneoppimismallin luominen Apache Spark MLlib:n avulla

Tässä artikkelissa opit käyttämään Apache Sparkia MLlib luodaksesi koneoppimissovelluksen, joka käsittelee ennakoivaa analyysiä Azure avoimella aineistolla. Spark tarjoaa sisäisiä koneoppimiskirjastoja. Tässä esimerkissä luokitusta käytetään logistista regressiota käyttämällä.

Tässä opetusohjelmassa käsitellään seuraavat vaiheet:

  • Aseta muistikirja ja tuonnit
  • Lataa ja näytteistä NYC:n taksidataa
  • Valmistele ja suunnittele ominaisuuksia
  • Koodaa kategorisia piirteitä
  • Junalogistiikan regressiomalli
  • Arvioi ja visualisoi tuloksia

SparkML- ja MLlib Spark -ydinkirjastot tarjoavat monia apuohjelmia, joista on hyötyä koneoppimistehtävissä. Nämä apuohjelmat soveltuvat seuraaviin:

  • Luokitus
  • Klusterointi
  • Hypoteesitestaus ja mallitilastojen laskeminen
  • Regressio
  • SVD-hajotus ja pääosa-analyysi (PCA)
  • Aiheen mallinnus

Edellytykset

Tutustu luokitukseen ja logistiseen regressioon

Luokittelu, suosittu koneoppimistehtävä, sisältää syötetietojen lajittelemisen luokkiin. Luokittelualgoritmi selvittää, miten tunnisteet annetaan syötetylle datalle. Esimerkiksi koneoppimisalgoritmi voisi hyväksyä osaketiedot syötteenä ja jakaa osakkeen kahteen kategoriaan: osakkeisiin, joita pitäisi myydä, ja osakkeisiin, jotka kannattaa säilyttää.

Logistinen regressioalgoritmi on hyödyllinen luokittelussa. Spark-logistinen regressio-ohjelmointirajapinta on hyödyllinen syötetietojen binaariluokittelussa yhteen kahdesta ryhmästä. Lisätietoja logistisesta regressiosta on Wikipediassa.

Logistinen regressio tuottaa logistisen funktion , joka ennustaa todennäköisyyden, että syötevektori kuuluu jompaan tai toiseen ryhmään.

Ennakoiva analyysi esimerkki New Yorkin kaupungin taksitiedoista

Tiedot ovat käytettävissä Azuren avoimet tietoaineistot -resurssin kautta. Tässä tietojoukon alijoukossa isännöidä tietoja keltaisista taksimatkoista, mukaan lukien alkamisajat, päättymisajat, alkamissijainnit, päättymissijainnit, matkakustannukset ja muut määritteet.

Tässä opetusohjelmassa käytetään Apache Sparkia analysoidakseen New Yorkin taksimatkan tippitietoja ja kehittääkseen mallin, jolla ennustataan, sisältääkö tietty matka tipin.

Luo Apache Spark -koneoppimismalli

  1. Luo PySpark-muistikirja. Lisätietoja löytyy kohdasta Luo muistikirja.

    Kun olet luonut muistikirjan, liitä se järventaloon valitsemalla vasemmasta paneelista Lisää järvitalo .

  2. Tuo tarvittavat tyypit tälle muistikirjalle. Liitä seuraava koodi ensimmäiseen soluun ja suorita se.

    import matplotlib.pyplot as plt
    from pyspark.sql.functions import unix_timestamp, date_format, col, when
    from pyspark.ml import Pipeline
    from pyspark.ml.feature import RFormula
    from pyspark.ml.feature import OneHotEncoder, StringIndexer
    from pyspark.ml.classification import LogisticRegression
    from pyspark.ml.evaluation import BinaryClassificationEvaluator
    

    Vahvista: Solu valmistuu ilman ImportError. Jos huomaat virheen, varmista, että kannettavasi käyttää PySpark-suoritusaikaa.

  3. Käytä MLflow'ta koneoppimiskokeiden ja niihin liittyvien ajojen seuraamiseen. Jos Microsoft Fabricin automaattinen loggaus on käytössä, vastaavat mittarit ja parametrit siepataan automaattisesti.

    import mlflow
    

    Vahvista: Solu valmistuu virheettömästi. Varmista, print(mlflow.__version__) että MLflow on saatavilla.

Muodosta syötteen DataFrame

Tämä esimerkki lataa Azuren avoimet tietoaineistot -tallennustilan tiedot Apache Spark DataFrameen. Sen jälkeen käytät Spark-operaatioita datan puhdistamiseen ja suodattamiseen.

  1. Liitä seuraava koodi uuteen soluun ja suorita se luodaksesi Spark DataFramen. Tämä vaihe hakee New Yorkin keltaisen taksin dataa, joka on suodatettu toukokuuhun 2018 asti.

    blob_account_name = "azureopendatastorage"
    blob_container_name = "nyctlc"
    blob_relative_path = "yellow"
    wasbs_path = f"wasbs://{blob_container_name}@{blob_account_name}.blob.core.windows.net/{blob_relative_path}"
    
    nyc_tlc_df = spark.read.parquet(wasbs_path) \
        .filter((col("tpepPickupDateTime") >= "2018-05-01") & (col("tpepPickupDateTime") < "2018-06-01")) \
        .repartition(20)
    

    Vahvista: Suorita seuraava solu varmistaaksesi datan latauksen onnistuneesti.

    print(f"Loaded {nyc_tlc_df.count()} rows")
    # Expected output: Loaded approximately 9,000,000+ rows
    
  2. Ota näyte aineistosta nopeuttaaksesi kehitystä ja koulutusta.

    # Sample without replacement to avoid duplicates
    sampled_taxi_df = nyc_tlc_df.sample(False, 0.001, seed=1234)
    

    Varmista: Varmista, että otoskoko on hallittavissa.

    print(f"Sampled {sampled_taxi_df.count()} rows")
    # Expected output: Sampled approximately 9,000-10,000 rows
    
  3. Katso display() dataa käyttämällä sisäänrakennettua komentoa tutkiaksesi datanäytettä.

    display(sampled_taxi_df.limit(10))
    

    Varmista: Tauluko, jossa on 10 riviä ja jossa näkyy sarakkeet kuten tpepPickupDateTime, fareAmount, tipAmount, ja tripDistance.

Tietojen valmistelu

Tietojen valmistelu on tärkeä vaihe koneoppimisprosessissa. Se sisältää raakadatan puhdistamisen, muuntamisen ja järjestämisen, jotta se soveltuu analyysiin ja mallinnukseen. Tässä osiossa suorita useita datan valmisteluvaiheita:

  • Suodata aineisto poistaaksesi poikkeavat arvot ja virheelliset arvot.
  • Poista sarakkeet, joita ei tarvita mallin koulutukseen.
  • Luo uudet sarakkeet raakadatasta.
  • Luo etiketti, jolla selvitetään, sisältyykö tiettyyn taksimatkaan tippi.

Suorita seuraava koodi valitaksesi relevantit sarakkeet, laskeaksesi johdetut ominaisuudet ja suodattaaksesi poikkeavat arvot:

taxi_df = sampled_taxi_df.select('totalAmount', 'fareAmount', 'tipAmount', 'paymentType', 'rateCodeId', 'passengerCount',
                    'tripDistance', 'tpepPickupDateTime', 'tpepDropoffDateTime',
                    date_format('tpepPickupDateTime', 'HH').cast('integer').alias('pickupHour'),
                    date_format('tpepPickupDateTime', 'EEEE').alias('weekdayString'),
                    (unix_timestamp(col('tpepDropoffDateTime')) - unix_timestamp(col('tpepPickupDateTime'))).alias('tripTimeSecs'),
                    (when(col('tipAmount') > 0, 1).otherwise(0)).alias('tipped')
                    ) \
            .filter((sampled_taxi_df.passengerCount > 0) & (sampled_taxi_df.passengerCount < 8)
                    & (sampled_taxi_df.tipAmount >= 0) & (sampled_taxi_df.tipAmount <= 25)
                    & (sampled_taxi_df.fareAmount >= 1) & (sampled_taxi_df.fareAmount <= 250)
                    & (sampled_taxi_df.tipAmount < sampled_taxi_df.fareAmount)
                    & (sampled_taxi_df.tripDistance > 0) & (sampled_taxi_df.tripDistance <= 100)
                    & (sampled_taxi_df.rateCodeId <= 5)
                    & (sampled_taxi_df.paymentType.isin({"1", "2"}))
                    )

Important

Funktio date_format käyttää mallia 'HH' (24 tunnin muoto, arvot 0–23) eikä 'hh' (12 tunnin muoto, arvot 1–12). 24 tunnin formaatti vaaditaan seuraavaan vuorokaudenaikaan perustuvan binning-logiikan vuoksi.

Seuraavaksi lisää liikenneaikabins-ominaisuus vuorokaudenajan mukaan:

taxi_featurised_df = taxi_df.select('totalAmount', 'fareAmount', 'tipAmount', 'paymentType', 'passengerCount',
                                    'tripDistance', 'weekdayString', 'pickupHour', 'tripTimeSecs', 'tipped',
                                    when((col('pickupHour') <= 6) | (col('pickupHour') >= 20), "Night")
                                    .when((col('pickupHour') >= 7) & (col('pickupHour') <= 10), "AMRush")
                                    .when((col('pickupHour') >= 11) & (col('pickupHour') <= 15), "Afternoon")
                                    .when((col('pickupHour') >= 16) & (col('pickupHour') <= 19), "PMRush")
                                    .otherwise("Other").alias('trafficTimeBins')
                                    ) \
                            .filter((taxi_df.tripTimeSecs >= 30) & (taxi_df.tripTimeSecs <= 7200))

Varmista: Varmista, että liikenneaikalaatikot on jaettu oikein.

taxi_featurised_df.groupBy('trafficTimeBins').count().show()
# Expected output: Shows counts for Night, AMRush, Afternoon, PMRush categories

Logistisen regressiomallin luominen

Viimeisessä tehtävässä otsikoidut tiedot muunnetaan muotoon, jota logistinen regressio pystyy käsittelemään. Logistisen regressioalgoritmin syötteellä on oltava otsikko-/ominaisuusvektoriparien rakenne, jossa ominaisuusvektori on syötepistettä edustavien lukujen vektori.

Muunna kategoriset sarakkeet trafficTimeBins ja weekdayString kokonaislukuesityksiksi käyttämällä seuraavaa OneHotEncoder lähestymistapaa:

# Convert categorical features into numeric representations
sI1 = StringIndexer(inputCol="trafficTimeBins", outputCol="trafficTimeBinsIndex")
en1 = OneHotEncoder(inputCol="trafficTimeBinsIndex", outputCol="trafficTimeBinsVec")
sI2 = StringIndexer(inputCol="weekdayString", outputCol="weekdayIndex")
en2 = OneHotEncoder(inputCol="weekdayIndex", outputCol="weekdayVec")

# Apply the encodings to create a new DataFrame
encoded_final_df = Pipeline(stages=[sI1, en1, sI2, en2]).fit(taxi_featurised_df).transform(taxi_featurised_df)

Vahvista: Vahvista, että koodattu DataFrame sisältää odotetut uudet sarakkeet.

print("Columns:", encoded_final_df.columns)
print(f"Row count: {encoded_final_df.count()}")
# Expected output: Columns list includes 'trafficTimeBinsVec' and 'weekdayVec'

Logistisen regressiomallin harjoittaminen

Jaa aineisto harjoitusjoukkoon (70%) ja testijoukkoon (30%):

# Split the DataFrame into training and test sets
trainingFraction = 0.7
testingFraction = (1 - trainingFraction)
seed = 1234

train_data_df, test_data_df = encoded_final_df.randomSplit([trainingFraction, testingFraction], seed=seed)

Varmista: Varmista, että jako tuotti kohtuulliset kokoiset.

print(f"Training rows: {train_data_df.count()}, Test rows: {test_data_df.count()}")
# Expected output: Approximately 70%/30% split of the encoded data

Luo mallikaava, kouluta logistinen regressiomalli ja arvioi se käyttämällä ROC:n (vastaanottajan toimintakyky) -käyrän pinta-alaa:

# Create a logistic regression model
logReg = LogisticRegression(maxIter=10, regParam=0.3, labelCol='label')

# Define the formula: 'tipped' is the response variable, right-hand side are predictors
classFormula = RFormula(formula="tipped ~ pickupHour + weekdayVec + passengerCount + tripTimeSecs + tripDistance + fareAmount + paymentType + trafficTimeBinsVec")

# Train the model using a pipeline
lrModel = Pipeline(stages=[classFormula, logReg]).fit(train_data_df)

# Generate predictions on the test dataset
predictions = lrModel.transform(test_data_df)

# Evaluate using Area Under ROC
evaluator = BinaryClassificationEvaluator(rawPredictionCol="rawPrediction", metricName="areaUnderROC")
auc = evaluator.evaluate(predictions)
print(f"Area under ROC = {auc}")

Vahvista: Tulos näyttää AUC-arvon. Hyvin toimiva malli tuottaa arvon lähellä 1,0.

Area under ROC = 0.97 (approximately)

Note

Tarkka AUC-arvo vaihtelee aineiston mukaan. Arvot yli 0,90 osoittavat vahvaa ennustesuorituskykyä tälle aineistolle.

Visuaalisen esityksen luominen ennusteesta

Rakenna lopullinen visualisointi tulkintaa varten mallin tuloksia. ROC-käyrä esittää todellisen positiivisen ja väärän positiivisen koron välisen kompromissin.

# Plot the ROC curve from the model training summary
modelSummary = lrModel.stages[-1].summary

# Extract FPR and TPR values as plain lists
roc_data = modelSummary.roc.select('FPR', 'TPR').toPandas()

plt.figure(figsize=(8, 6))
plt.plot([0, 1], [0, 1], 'r--', label='Random classifier')
plt.plot(roc_data['FPR'], roc_data['TPR'], label=f'Logistic Regression (AUC = {auc:.4f})')
plt.xlabel('False Positive Rate')
plt.ylabel('True Positive Rate')
plt.title('ROC Curve - NYC Taxi Tip Prediction')
plt.legend(loc='lower right')
plt.show()

Vahvista: Kuvaaja näyttää ROC-käyrän punaisen katkoviivan yläpuolella. Käyrän tulisi kaartua kohti vasempaan yläkulmaan, mikä osoittaa vahvaa luokitussuoritusta.

Kaavio, joka näyttää logistista regressiota vastaavan ROC-käyrän kärkimallissa.

Puhdista resurssit

Kun olet suorittanut tämän tutoriaalin, poista muistikirja ja järvitalo vapauttaaksesi työtilan kapasiteettia:

  1. Työtilassasi napsauta muistikirjaa hiiren oikealla ja valitse Poista.
  2. Jos loit järvenrakennuksen nimenomaan tätä opetusta varten, napsauta sitä hiiren oikealla ja valitse Poista.

Koulutetun mallin säilyttämiseksi tulevaa käyttöä varten lisää seuraava koodi ennen puhdistusta:

# Save the model to the lakehouse
model_path = "abfss://<your-workspace>@onelake.dfs.fabric.microsoft.com/<your-lakehouse>.Lakehouse/Files/models/taxi_tip_model"
lrModel.write().overwrite().save(model_path)
print(f"Model saved to: {model_path}")

Vianmääritys

Ongelma Syy Ratkaisu
Py4JJavaError Parquetia lukiessa Verkkoyhteys Azure-blob-tallennukseen Varmista, että Fabric-työtilassasi on ulospäin suuntautuva internet-yhteys. Kokeile käynnistää Spark-istunto uudelleen.
AnalysisException: cannot resolve column Sarakkeen nimen kirjoitusvirhe tai skeeman ristiriita Juokse nyc_tlc_df.printSchema() tarkastamaan saatavilla olevat sarakkeet. NYC:n taksiaineiston skeema voi muuttua vuosien välillä.
Tyhjennä datakehys suodatuksen jälkeen Suodatinehdot ovat liian rajoittavia data-ikkunalle Laajenna päivämääräväliä tai tarkista sampled_taxi_df.count() ennen suodatusta.
IllegalArgumentException StringIndexerissä Näkymättömät etiketit muunnoksen aikana Lisää handleInvalid="skip" puheluihisi StringIndexer : StringIndexer(inputCol="...", outputCol="...", handleInvalid="skip")
Alhainen AUC (alle 0,6) Riittämätön data tai virheellinen ominaisuussuunnittelu Nosta otososuutta (esimerkiksi 0.01 sen sijaan 0.001) ja varmista trafficTimeBins , että kategoriat ovat tasapainossa.
OutOfMemoryError Aineisto liian suuri käytettävissä olevaan kapasiteettiin nähden Vähennä näytteen osuutta tai kasvata Fabric-kapasiteettitasoa.
ROC-kuvaaja ei näy Matplotlib-taustajärjestelmä muistikirjassa Lisää %matplotlib inline vihkon yläosaan.