Apache Spark MLlib ile makine öğrenmesi modeli oluşturma

Bu makalede apache Spark MLlib kullanarak açık Azure bir veri kümesinde tahmine dayalı analizi işleyen bir makine öğrenmesi uygulaması oluşturmayı öğreneceksiniz. Spark, yerleşik makine öğrenmesi kitaplıkları sağlar. Bu örnek lojistik regresyon aracılığıyla sınıflandırmayı kullanır.

Bu öğretici şu adımları kapsar:

  • Not defterini ve içeri aktarmaları ayarlama
  • NYC taksi verilerini yükleme ve örnekleme
  • Hazırlama ve mühendislik özellikleri
  • Kategorik özellikleri kodlama
  • Lojistik regresyon modelini eğitme
  • Sonuçları değerlendirme ve görselleştirme

Temel SparkML ve MLlib Spark kitaplıkları, makine öğrenmesi görevleri için yararlı olan birçok yardımcı program sağlar. Bu yardımcı programlar şunlar için uygundur:

  • Sınıflandırma
  • Kümeleme
  • Hipotez testi ve örnek istatistikleri hesaplama
  • Regresyon
  • Tekil değer ayrıştırma (SVD) ve asıl bileşen analizi (PCA)
  • Konu modelleme

Prerequisites

Sınıflandırmayı ve lojistik regresyonu anlama

Popüler bir makine öğrenmesi görevi olan sınıflandırma, giriş verilerini kategorilere ayırmayı içerir. Sınıflandırma algoritması, sağlanan giriş verilerine etiketlerin nasıl atandığını açıklar. Örneğin, bir makine öğrenmesi algoritması hisse senedi bilgilerini giriş olarak kabul edebilir ve hisse senedini iki kategoriye bölebilir: satmanız gereken hisse senetleri ve tutmanız gereken hisse senetleri.

Lojistik regresyon algoritması sınıflandırma için kullanışlıdır. Spark lojistik regresyon API'si, giriş verilerinin iki gruptan birinde ikili sınıflandırması için kullanışlıdır. Lojistik regresyon hakkında daha fazla bilgi için bkz . Vikipedi.

Lojistik regresyon, giriş vektörlerinin bir gruba veya diğerine ait olma olasılığını tahmin eden bir lojistik işlev oluşturur.

NYC taksi verilerinin tahmine dayalı analiz örneği

Veriler Azure Açık Veri Kümeleri kaynağı üzerinden kullanılabilir. Bu veri kümesi, başlangıç saatleri, bitiş saatleri, başlangıç konumları, bitiş konumları, seyahat maliyetleri ve diğer öznitelikler de dahil olmak üzere sarı taksi yolculukları hakkındaki bilgileri barındırır.

Bu öğreticide, NYC taksi yolculuklarına ait bahşiş verileri üzerinde analiz yapmak ve belirli bir yolculukta bahşiş verilip verilmediğini tahmin eden bir model geliştirmek için Apache Spark kullanılır.

Apache Spark makine öğrenmesi modeli oluşturma

  1. PySpark not defteri oluşturun. Daha fazla bilgi için bkz. Not defteri oluşturma.

    Not defterini oluşturduktan sonra, sol bölmede Lakehouse ekle'yi seçerek onu bir lakehouse'a bağlayın.

  2. Bu not defteri için gerekli türleri içeri aktarın. Aşağıdaki kodu ilk hücreye yapıştırın ve çalıştırın.

    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
    

    Doğrula: Hücre, ImportError olmadan tamamlanır. Hata görürseniz not defterinizin PySpark çalışma zamanını kullandığını onaylayın.

  3. Makine öğrenmesi denemelerinizi ve buna karşılık gelen çalıştırmaları izlemek için MLflow kullanın. Microsoft Fabric Otomatik Kaydetme etkinleştirildiyse, ilgili ölçümler ve parametreler otomatik olarak yakalanır.

    import mlflow
    

    Doğrula: Hücre hatasız tamamlar. MLflow'un kullanılabilir olduğunu onaylamak için komutunu çalıştırın print(mlflow.__version__) .

Giriş DataFrame'ini oluşturma

Bu örnek, Azure Açık Veri Kümeleri depolamadaki verileri bir Apache Spark DataFrame'e yükler. Ardından, veri kümesini temizlemek ve filtrelemek için Spark işlemlerini uygularsınız.

  1. Aşağıdaki kodu yeni bir hücreye yapıştırın ve Spark DataFrame oluşturmak için çalıştırın. Bu adım, Mayıs 2018'e filtrelenen NYC sarı taksi verilerini alır.

    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)
    

    Doğrula: Verilerin başarıyla yüklendiğini onaylamak için aşağıdaki hücreyi çalıştırın.

    print(f"Loaded {nyc_tlc_df.count()} rows")
    # Expected output: Loaded approximately 9,000,000+ rows
    
  2. Geliştirme ve eğitimi hızlandırmak için veri kümesini örnekleyin.

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

    Doğrula: Örnek boyutunun yönetilebilir olduğunu onaylayın.

    print(f"Sampled {sampled_taxi_df.count()} rows")
    # Expected output: Sampled approximately 9,000-10,000 rows
    
  3. Veri örneğini keşfetmek için yerleşik display() komutunu kullanarak verileri görüntüleyin.

    display(sampled_taxi_df.limit(10))
    

    Doğrula: tpepPickupDateTime, fareAmount, tipAmount ve tripDistance gibi sütunların gösterildiği 10 satırlı bir tablo görüntülenir.

Verileri hazırlama

Veri hazırlama, makine öğrenmesi sürecinde önemli bir adımdır. Analiz ve modelleme için uygun hale getirmek için ham verilerin temizlenmesini, dönüştürülmesini ve düzenlenmesini içerir. Bu bölümde, çeşitli veri hazırlama adımları gerçekleştirin:

  • Aykırı değerleri ve yanlış değerleri kaldırmak için veri kümesini filtreleyin.
  • Model eğitimi için gerekli olmayan sütunları kaldırın.
  • Ham verilerden yeni sütunlar oluşturun.
  • Verilen bir taksi yolculuğunun bahşiş içerip içermediğini belirlemek için bir etiket oluşturun.

İlgili sütunları seçmek, türetilmiş özellikleri hesaplamak ve aykırı değerleri filtrelemek için aşağıdaki kodu çalıştırın:

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

İşlev, date_format (12 saatlik biçim, 1-12 değerleri) yerine 'HH' deseni 'hh' (24 saatlik biçim, değerler 0-23) kullanır. 24 saatlik biçim, izleyen günün saati gruplama mantığı için gereklidir.

Ardından, günün saatlerine göre trafik zaman bölmeleri özelliğini ekleyin:

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

Doğrula: Trafik zaman bölmelerinin doğru dağıtıldığını onaylayın.

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

Lojistik regresyon modeli oluşturma

Son görev etiketlenmiş verileri lojistik regresyonun işleyebileceği bir biçime dönüştürür. Lojistik regresyon algoritması için giriş, özellik vektörü giriş noktasını temsil eden sayıları içeren bir vektör olan etiket/özellik vektör çiftleri yapısında olmalıdır.

Yaklaşımı kullanarak kategorik sütunları trafficTimeBins ve weekdayString tamsayı gösterimlerine dönüştürün 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)

Doğrula: Kodlanmış DataFrame'in beklenen yeni sütunlara sahip olduğunu onaylayın.

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

Lojistik regresyon modelini eğitin

Veri kümesini bir eğitim kümesine (70%) ve bir test kümesine (30%) bölün:

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

Doğrulayın: Bölme işleminin makul boyutlar ürettiğini onaylayın.

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

Model formülünü oluşturun, lojistik regresyon modelini eğitin ve ROC (Alıcı Çalışma Özelliği) Eğrisi Altındaki Alan'ı kullanarak değerlendirin:

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

Doğrulama: Çıktıda bir AUC değeri gösterilir. İyi performans gösteren bir model 1,0'a yakın bir değer üretir.

Area under ROC = 0.97 (approximately)

Note

Tam AUC değeri, veri örneğine bağlı olarak değişir. 0,90'ın üzerindeki değerler bu veri kümesi için güçlü tahmine dayalı performans gösterir.

Tahminin görsel gösterimini oluşturma

Model sonuçlarını yorumlamak için son bir görselleştirme oluşturun. ROC eğrisi, gerçek pozitif oran ile hatalı pozitif oran arasındaki dengeyi gösterir.

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

Doğrula: Kırmızı kesikli çapraz çizginin üzerinde ROC eğrisini gösteren bir çizim görüntülenir. Eğri, güçlü sınıflandırma performansını gösteren sol üst köşeye doğru eğilmelidir.

Tip modelinde lojistik regresyon için ROC eğrisini gösteren grafik.

Kaynakları temizle

Bu öğreticiyi tamamladıktan sonra, çalışma alanı kapasitesini boşaltmak için not defterini ve lakehouse'u silin:

  1. Çalışma alanınızda not defterine sağ tıklayın ve Sil'i seçin.
  2. Bu öğreticiye özel olarak bir lakehouse oluşturduysanız, üzerine sağ tıklayın ve Sil'i seçin.

Gelecekte kullanmak üzere eğitilmiş modeli korumak için, temizleme işleminden önce aşağıdaki kodu ekleyin:

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

Sorun giderme

Sorun Nedeni Çözüm
Py4JJavaError parquet okurken Azure Blob Depolama'ya ağ bağlantısı Fabric çalışma alanınızın giden İnternet erişimi olduğunu doğrulayın. Spark oturumunu yeniden başlatmayı deneyin.
AnalysisException: cannot resolve column Sütun adı yazım hatası veya şema uyuşmazlığı Kullanılabilir sütunları incelemek için komutunu çalıştırın nyc_tlc_df.printSchema() . NYC taksi veri kümesi şeması yıllar arasında değişebilir.
Filtrelemeden sonra DataFrame'i boşaltma Filtre koşulları veri penceresi için çok kısıtlayıcı Filtrelemeden önce tarih aralığını artırın veya denetleyin sampled_taxi_df.count() .
IllegalArgumentException StringIndexer'da Dönüştürme sırasında etiketlerin görünmemesi Aramalarınıza handleInvalid="skip" ekleyinStringIndexer:StringIndexer(inputCol="...", outputCol="...", handleInvalid="skip")
Düşük AUC (0,6'nın altında) Yetersiz veri veya yanlış özellik mühendisliği Örneklem oranını artırın (örneğin, 0.01 yerine 0.001) ve trafficTimeBins kategorilerinin dengeli olduğunu doğrulayın.
OutOfMemoryError Kullanılabilir kapasite için veri kümesi çok büyük Örnekleme oranını azaltın veya Fabric kapasite düzeyinizi artırın.
ROC çizimi görüntülenmiyor Not defterinde Matplotlib arka uç sorunu Not defterinin en üstüne ekleyin %matplotlib inline .