使用 Apache Spark MLlib 建置機器學習模型

在本文中,您將學習如何使用 Apache Spark MLlib 來建立一個機器學習應用程式,能在Azure開放資料集中進行預測分析。 Spark 提供內建的機器學習程式庫。 此範例會透過羅吉斯迴歸使用分類。

本教學課程涵蓋了下列步驟:

  • 設定筆記本並匯入項目
  • 載入並取樣紐約市計程車資料
  • 準備及設計功能
  • 編碼類別型特徵
  • 列車邏輯迴歸模型
  • 評估並視覺化結果

核心 SparkML 和 MLlib SparkMLlib 程式庫提供許多可用於機器學習工作的工具。 這些公用程式適用於:

  • 分類
  • 叢集
  • 假設測試和計算範例統計資料
  • 迴歸
  • 奇異值分解 (SVD) 和主體元件分析 (PCA)
  • 主題模型化

先決條件

了解分類和邏輯迴歸

分類是常見的機器學習工作,涉及將輸入資料依類別排序。 分類演算法會找出如何為所提供輸入資料指派 標籤 的方法。 例如,機器學習演算法可以接受股票資訊作為輸入,並將股票分為兩類:應該賣出的股票和應該保留的股票。

邏輯斯迴歸演算法對分類非常有用。 Spark 的羅吉斯迴歸 API 可用於二元分類,將輸入資料歸類到兩個群組之一。 如需有關羅吉斯迴歸的詳細資訊,請參閱 Wikipedia。

邏輯斯迴歸會產生一個 logistic 函數,用來預測輸入向量屬於兩個群組其中一組的機率。

紐約市計程車資料的預測性分析範例

此資料可透過 Azure 開放資料集資源取得。 此資料集子集裝載黃色計程車車程的相關資訊,包括開始時間、結束時間、開始位置、結束位置、車程成本和其他屬性。

本教學使用 Apache Spark 分析紐約市計程車小費資料,並建立模型來預測特定行程是否包含小費。

建立 Apache Spark 機器學習模型

  1. 建立 PySpark 筆記本。 欲了解更多資訊,請參閱 建立筆記本。

    建立筆記本後,在左側面板選擇 「新增湖屋 」,將其附加到湖屋上。

  2. 匯入此筆記本所需的型別。 把以下程式碼貼到第一個儲存格並執行。

    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
    

    驗證:該單元在沒有 ImportError的情況下完備。 如果你看到錯誤,請確認你的筆記本是否使用 PySpark 執行時。

  3. 使用 MLflow 來追蹤你的機器學習實驗和相應的執行。 如果已啟用 Microsoft Fabric 自動記錄,系統會自動擷取對應的計量和參數。

    import mlflow
    

    驗證:儲存格完整且無錯誤。 執行 print(mlflow.__version__) 以確認 MLflow 可用。

建立輸入的 DataFrame

此範例將 Azure 開放資料集 儲存的資料載入 Apache Spark DataFrame。 接著,你要用 Spark 操作來清理和篩選資料集。

  1. 將以下程式碼貼到新儲存格,執行以建立 Spark DataFrame。 此步驟可取得篩選至2018年5月的紐約市黃色計程車資料。

    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)
    

    驗證:執行以下儲存格以確認資料載入成功。

    print(f"Loaded {nyc_tlc_df.count()} rows")
    # Expected output: Loaded approximately 9,000,000+ rows
    
  2. 取樣資料集以加快開發與訓練。

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

    驗證:確認樣本量是否可管理。

    print(f"Sampled {sampled_taxi_df.count()} rows")
    # Expected output: Sampled approximately 9,000-10,000 rows
    
  3. 使用內建 display() 指令瀏覽資料樣本。

    display(sampled_taxi_df.limit(10))
    

    驗證:會出現一個有 10 列的表格,顯示欄位如 tpepPickupDateTime、 fareAmount、 tipAmounttripDistance、 和 。

準備資料

資料準備是機器學習程序中的重要步驟。 它涉及清理、轉換並組織原始資料,使其適合分析與建模。 在本節中,執行數個資料準備步驟:

  • 過濾資料集以移除異常值和錯誤值。
  • 移除不需要用於模型訓練的欄位。
  • 從原始資料建立新的欄位。
  • 產生標籤來判斷某趟計程車是否包含小費。

執行以下程式碼以選擇相關欄位、計算衍生特徵,並過濾離群值:

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

該 date_format 函數使用模式 'HH' (24小時格式,值0-23)而非 'hh' (12小時格式,值1-12)。 24 小時格式是後續時間分區邏輯所必需的。

接著,根據一天中的時段加入交通時間箱功能:

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

驗證:確認流量時間區塊分配正確。

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

建立羅吉斯迴歸模型

最終工作會將標示的資料轉換成羅吉斯迴歸可以處理的格式。 羅吉斯迴歸演算法的輸入必須具有標籤/特徵向量組結構,其中特徵向量是代表輸入點的數字向量。

將類別欄位 trafficTimeBins 和 weekdayString 轉換為整數表示,方法是: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)

驗證:確認編碼後的資料框架是否具備預期的新欄位。

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

訓練羅吉斯迴歸模型

將資料集拆分為訓練集(70 個%)和一個測試集(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)

確認:確認分拆產生的尺寸是否合理。

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

建立模型公式,訓練邏輯迴歸模型,並利用 ROC(接收者操作特徵)曲線下的面積來評估:

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

驗證:輸出顯示 AUC 值。 一個表現良好的模型會產生接近 1.0 的數值。

Area under ROC = 0.97 (approximately)

Note

具體的AUC值會依資料樣本而異。 值高於 0.90 表示該資料集的預測表現良好。

建立預測的視覺表示法

建立最終的視覺化以解讀模型結果。 ROC曲線呈現真陽性率與假陽性率之間的權衡。

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

驗證:顯示紅色虛線對角線上方的 ROC 曲線圖。 曲線應向左上角彎曲,顯示出優異的分類表現。

顯示小費模型中羅吉斯迴歸之 ROC 曲線的圖表。

清理資源

完成這個教學後,刪除筆記本和 Lakehouse 以釋放工作空間容量:

  1. 在你的工作區裡,右鍵點選筆記本並選擇 刪除。
  2. 如果你是特地為本教學建立此湖屋,請以滑鼠右鍵按一下它,然後選取 刪除。

為了保留訓練好的模型以備未來使用,請在清理前加入以下程式碼:

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

Troubleshooting

Issue 原因 解決方法
Py4JJavaError 讀 Parquet 時 Azure Blob 儲存體的網路連線 確認你的 Fabric 工作區有外接網路。 試著重新啟動 Spark session。
AnalysisException: cannot resolve column 欄位名稱的錯字或結構不符 跑 nyc_tlc_df.printSchema() 去檢查可用的欄位。 紐約市計程車資料集架構可能在不同年份之間有所變動。
篩選後的空資料框 過濾條件對資料視窗來說過於限制 在篩選前請擴大日期範圍或檢查 sampled_taxi_df.count() 。
IllegalArgumentException 位於 StringIndexer 中 轉換過程中看不見的標籤 將 handleInvalid="skip" 加入您的 StringIndexer 呼叫中:StringIndexer(inputCol="...", outputCol="...", handleInvalid="skip")
低AUC(低於0.6) 資料不足或特徵工程錯誤 提高樣本比例(例如, 0.01 取代 0.001),並驗證 trafficTimeBins 類別是否平衡。
OutOfMemoryError 資料集容量過大,無法滿足可用容量 減少樣本比例或提升你的 Fabric 容量等級。
ROC 圖未顯示 筆記本中的 Matplotlib 後端問題 在筆記本頂端加上 %matplotlib inline 。