Bygg en maskinlæringsmodell med Apache Spark MLlib

I denne artikkelen lærer du hvordan du bruker Apache Spark MLlib for å lage en maskinlæringsapplikasjon som håndterer prediktiv analyse på et Azure åpent datasett. Spark gir innebygde maskinlæringsbiblioteker. Dette eksemplet bruker klassifisering gjennom logistisk regresjon.

Denne opplæringen dekker disse trinnene:

  • Oppsett notatbok og import
  • Last inn og prøv NYC-taxidata
  • Forbered og konstruer funksjoner
  • Kod kategoriske egenskaper
  • Toglogistisk regresjonsmodell
  • Evaluer og visualiser resultatene

Kjernebibliotekene SparkML og MLlib Spark gir mange verktøy som er nyttige for maskinlæringsoppgaver. Disse verktøyene er egnet for:

  • Klassifisering
  • Klynging
  • Hypotesetesting og beregning av eksempelstatistikk
  • Regresjon
  • Entallsverdikomponering (SVD) og hovedkomponentanalyse (PCA)
  • Emnemodellering

Forutsetninger

  • Få et Microsoft Fabric-abonnement. Eller registrer deg for en gratis prøveversjon av Microsoft Fabric.

  • Logg på Microsoft Fabric.

  • Bytt til Fabric ved å bruke erfaringsbryteren nederst til venstre på hjemmesiden din.

    Skjermbilde som viser valget av Fabric i menyen for opplevelsesbytter.

Forstå klassifisering og logistisk regresjon

Klassifisering, en populær maskinlæringsoppgave, innebærer sortering av inndata i kategorier. En klassifiseringsalgoritme finner ut hvordan den skal tildele etiketter til de oppgitte inputdataene. For eksempel kan en maskinlæringsalgoritme akseptere aksjeinformasjon som input og dele aksjen inn i to kategorier: aksjer du bør selge og aksjer du bør beholde.

Den logistiske regresjonsalgoritmen er nyttig for klassifisering. Spark-logistikkregresjons-API-en er nyttig for binær klassifisering av inndata i én av to grupper. Hvis du vil ha mer informasjon om logistisk regresjon, kan du se Wikipedia.

Logistisk regresjon produserer en logistisk funksjon som forutsier sannsynligheten for at en inputvektor tilhører den ene eller den andre gruppen.

Eksempel på prediktiv analyse av NYC-taxidata

Dataene er tilgjengelige via Azure Open Datasets-ressursen. Dette delsettet for datasettet inneholder informasjon om gule taxiturer, inkludert starttider, sluttidspunkt, startsteder, sluttsteder, reisekostnader og andre attributter.

Denne veiledningen bruker Apache Spark til å analysere NYC-taxiturens tipsdata og utvikle en modell for å forutsi om en bestemt tur inkluderer et tips.

Opprette en Apache Spark-maskinlæringsmodell

  1. Opprett en PySpark-notatblokk. For mer informasjon, se Lag en notatbok.

    Etter at du har opprettet notatboken, fester du den til et lakehouse ved å velge Add lakehouse i venstre panel.

  2. Importer de nødvendige typene for denne notatboken. Lim inn følgende kode i den første cellen og kjør den.

    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
    

    Verifiser: Cellen fullføres uten ImportError. Hvis du ser en feil, bekreft at notatboken din bruker PySpark-runtime.

  3. Bruk MLflow til å spore maskinlæringseksperimentene dine og tilhørende kjøringer. Hvis Microsoft Fabric Autologging er aktivert, registreres de tilsvarende måledataene og parameterne automatisk.

    import mlflow
    

    Verifiser: Cellen fullføres uten feil. Kjør print(mlflow.__version__) for å bekrefte at MLflow er tilgjengelig.

Konstruere datarammen for inndata

Dette eksempelet laster data fra Azure Open Datasets Storage inn i en Apache Spark DataFrame. Deretter bruker du Spark-operasjoner for å rense og filtrere datasettet.

  1. Lim inn følgende kode i en ny celle og kjør den for å lage en Spark DataFrame. Dette steget henter NYC gule taxidata filtrert til mai 2018.

    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)
    

    Bekreft: Kjør følgende celle for å bekrefte at data lastes vellykket.

    print(f"Loaded {nyc_tlc_df.count()} rows")
    # Expected output: Loaded approximately 9,000,000+ rows
    
  2. Ta et utvalg av datasettet for å fremskynde utvikling og opplæring.

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

    Bekreft: Bekreft at utvalget er håndterbart.

    print(f"Sampled {sampled_taxi_df.count()} rows")
    # Expected output: Sampled approximately 9,000-10,000 rows
    
  3. Se dataene ved å bruke den innebygde display() kommandoen for å utforske dataprøven.

    display(sampled_taxi_df.limit(10))
    

    Verifiser: En tabell med 10 rader vises som kolonner som tpepPickupDateTime, fareAmount, tipAmount, og tripDistance.

Klargjøre dataene

Dataforberedelse er et viktig trinn i maskinlæringsprosessen. Det innebærer å rense, transformere og organisere rådata for å gjøre dem egnet for analyse og modellering. I denne seksjonen utføres flere steg for dataforberedelse:

  • Filtrer datasett for å fjerne uteliggere og feilaktige verdier.
  • Fjern kolonner som ikke trengs for modelltrening.
  • Lag nye kolonner fra rådataene.
  • Lag en etikett for å avgjøre om en taxitur innebærer tips.

Kjør følgende kode for å velge relevante kolonner, beregne avledede egenskaper og filtrere uteliggere:

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

Funksjonen date_format bruker mønsteret 'HH' (24-timers format, verdier 0-23) i stedet for 'hh' (12-timers format, verdier 1-12). 24-timers formatet kreves for time-of-day binning-logikken som følger.

Deretter legger du til funksjonen trafikktidsbinger basert på timen på døgnet:

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

Verifiser: Bekreft at trafikktidsbinene er riktig fordelt.

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

Opprette en logistisk regresjonsmodell

Den endelige oppgaven konverterer de merkede dataene til et format som logistikkregresjon kan håndtere. Inndataene til en logistisk regresjonsalgoritme må ha en struktur for etikett/funksjonsvektorpar, der funksjonsvektoren er en vektor av tall som representerer inndatapunktet.

Konverter de kategoriske kolonnene trafficTimeBins og weekdayString til heltallsrepresentasjoner ved å bruke tilnærmingen OneHotEncoder :

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

Bekreft: Bekreft at den kodede DataFrame har de forventede nye kolonnene.

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

Lære opp en logistisk regresjonsmodell

Del datasettet i et treningssett (70%) og et testsett (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)

Bekreft: Bekreft at splittelsen ga rimelige størrelser.

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

Lag modellformelen, tren den logistiske regresjonsmodellen, og evaluer den ved å bruke areal under ROC (Receiver Operating Characteristic) Curve:

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

Verifiser: Utgangen viser en AUC-verdi. En godt presterende modell gir en verdi nær 1,0.

Area under ROC = 0.97 (approximately)

Note

Den eksakte AUC-verdien varierer avhengig av dataprøven. Verdier over 0,90 indikerer sterk prediktiv ytelse for dette datasettet.

Opprette en visuell representasjon av prognosen

Bygg en endelig visualisering for å tolke modellens resultater. En ROC-kurve presenterer avveiningen mellom sann positivrate og falsk positivrate.

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

Verifiser: En graf vises som viser ROC-kurven over den røde stiplede diagonallinjen. Kurven skal bøye seg mot øvre venstre hjørne, noe som indikerer sterk klassifiseringsprestasjon.

Graf som viser ROC-kurven for logistisk regresjon i tipsmodellen.

Rydd opp ressurser

Når du er ferdig med denne veiledningen, slett notatboken og lakehouse for å frigjøre arbeidsplass:

  1. I arbeidsområdet ditt, høyreklikk på notatboken og velg Slett.
  2. Hvis du har laget et innsjøhus spesielt for denne veiledningen, høyreklikk på det og velg Slett.

For å bevare den trente modellen for fremtidig bruk, legg til følgende kode før opprydding:

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

Feilsøking

Problem Årsak Løsning
Py4JJavaError Når man leser parkett Network connectivity to Azure blob storage Sjekk at Fabric-arbeidsplassen din har utgående internett-tilgang. Prøv å starte Spark-økten på nytt.
AnalysisException: cannot resolve column Stavefeil i kolonnenavn eller skjemamismatch Løp nyc_tlc_df.printSchema() for å inspisere tilgjengelige kolonner. NYC-taxidatasettskjemaet kan endre seg mellom årene.
Tom DataFrame etter filtrering Filterbetingelser for restriktive for datavinduet Øk datointervallet eller sjekk sampled_taxi_df.count() før du filtrerer.
IllegalArgumentException i StringIndexer Usynlige etiketter under transformasjon Legg til handleInvalid="skip" flere StringIndexer samtaler: StringIndexer(inputCol="...", outputCol="...", handleInvalid="skip")
Lav AUC (under 0,6) Utilstrekkelig data eller feil funksjonsutvikling Øk prøvefraksjonen (for eksempel 0.01 i stedet for 0.001) og verifiser trafficTimeBins at kategoriene er balanserte.
OutOfMemoryError Datasett er for stort for tilgjengelig kapasitet Reduser prøvefraksjonen eller øk Fabric-kapasitetsnivået ditt.
ROC-plott vises ikke Matplotlib backend-problem i notatboken Legg til %matplotlib inline øverst i notatboken.