本文說明如何用兩種方法執行文本分類任務。 一種方法是單純使用 pyspark,另一種方法則是使用 synapseml 函式庫。 兩種方法都能產生相同的效能,但也突顯出 SynapseML 相較於 pyspark 可降低程式碼複雜度。
該任務根據評論文字預測亞馬遜上銷售書籍的顧客評論是好(評分 > 3)還是差評。 你用不同的超參數訓練 LogisticRegression 學習者,然後選擇最佳模型。
必要條件
取得 Microsoft Fabric 訂用帳戶。 或註冊免費的 Microsoft Fabric 試用版。
登入 Microsoft Fabric。
使用首頁左下角的體驗切換器切換到 Fabric。
- 建立 筆記本。
- 將您的筆記本連接到 Lakehouse。 在筆記本中,選擇左側窗格的 新增 ,即可附加現有湖畔別墅或建立新湖畔別墅。
Note
本文中使用的所有函式庫(pyspark、synapseml、numpy)皆預裝於 Fabric Spark 執行時中。 你不需要安裝任何套件。
載入並探索資料
在Fabric筆記本中,Spark session 已經以 spark 變數形式存在。 從公開的 Azure Blob 儲存體 位置載入 Amazon Book Review 資料集:
rawData = spark.read.parquet(
"wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet"
)
rawData.show(5)
確認資料集是否正確載入:
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")
擷取特徵和處理資料
實物資料通常具有多種特徵,例如文字、數字和類別資料。 為了展示如何處理混合特徵類型,請在資料集中加入兩個數值特徵:評論的 字數 與 平均字長。
定義使用者自訂函數(UDF)
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())
使用 SynapseML UDFTransformer 套用 UDF
使用來自 SynapseML 的 UDFTransformer,將 UDF 包裝成與管線相容的轉換器:
from synapse.ml.stages import UDFTransformer
wordLengthTransformer = UDFTransformer(
inputCol="text", outputCol="wordLength", udf=wordLengthUDF
)
wordCountTransformer = UDFTransformer(
inputCol="text", outputCol="wordCount", udf=wordCountUDF
)
執行功能流程
套用這兩個轉換器,並根據評分建立一個二元標籤欄位:
from pyspark.ml import Pipeline
data = (
Pipeline(stages=[wordLengthTransformer, wordCountTransformer])
.fit(rawData)
.transform(rawData)
.withColumn("label", rawData["rating"] > 3)
.drop("rating")
)
驗證特徵提取:
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")
使用 pyspark 分類
要使用函式庫選擇最佳的 LogisticRegression 分類器 pyspark ,您必須明確執行以下步驟:
- 處理各項功能:
- 用詞彙化文字欄位。
- 利用雜湊將分詞化欄位雜湊成向量。
- 將數值特徵與向量合併。
- 將標籤欄位從布林類型轉換為整數類型。
- 在
train資料集上使用不同的超參數訓練多個 LogisticRegression 演算法。 - 計算每個訓練模型的 ROC 曲線下面積(AUC),並選擇資料集中
test指標最高的模型。 - 評估
validation集合上的最佳模型。
將資料特徵化並加以準備
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())
)
請驗證所附資料:
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")
訓練與評估模型
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}")
驗證結果:
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}")
Note
具體的AUC值取決於隨機分配。 預計值介於0.65到0.85之間。
使用 SynapseML 分類
這種 synapseml 方法以更少的步驟達成相同的效果。 SynapseML 內部處理特徵化,減少了你需要撰寫的程式碼:
-
TrainClassifier估計器會在內部將資料特徵化,前提是train、test和validation資料集中的欄位代表特徵。 -
FindBestModel估計器透過評估該資料集在指定指標上的test效能,從訓練模型池中找出最佳模型。 -
ComputeModelStatistics轉換器會同時針對已評分的資料集(在此情況下為validation資料集)計算多項指標。
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)
)
驗證SynapseML的結果:
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}")
Note
pyspark 與 SynapseML 方法應該會產生相似的 AUC 值,因為它們訓練相同的模型類型與相同的超參數,且資料相同。
比較這兩種方法
| Aspect | pyspark | SynapseML |
|---|---|---|
| 特徵處理 | 手冊(從 Tokenizer 到 HashingTF 再到 VectorAssembler) | 自動(由 TrainClassifier 處理) |
| 型號選擇 | 手動迴路與評估器 | 內建 FindBestModel |
| 度量計算 | 每次評估呼叫僅對應一個指標 | 多重指標 ComputeModelStatistics |
| 程式碼行數 | 大約30行 | 大約15行 |
| Result | 相同的 AUC | 相同的 AUC |
Troubleshooting
| Issue | 原因 | Resolution |
|---|---|---|
AnalysisException: Path does not exist |
公開的 blob 儲存 URL 暫時無法使用 | 請稍候幾分鐘再重試。 透過執行 spark.read.parquet("wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet").count() 來驗證連線能力 |
IllegalArgumentException: Field "features" does not exist |
變壓器之間的功能欄位名稱不一致 | 在 VectorAssembler 步驟之前,請先執行 data.columns 以驗證欄位名稱 |
NameError: name 'LogisticRegression' is not defined |
缺少匯入說明 | 在儲存格頂端加from pyspark.ml.classification import LogisticRegression |
ModuleNotFoundError: No module named 'synapse.ml' |
Notebook 未使用 Fabric Spark 執行階段 | 確認筆記本是否使用 Fabric Runtime 1.2 或更新版本。 在色帶中選擇 環境 來檢查。 |
| 低AUC(低於0.6) | 資料分割異常或收斂問題 | 使用 data.groupBy("label").count().show() 驗證標籤分布。 預期資料集大致平衡。 |
Py4JJavaError: An error occurred while calling |
Java/Spark 內部錯誤 | 請查看 Spark 介面上的詳細錯誤日誌。 選取 Session>停止工作階段 以重新啟動 Spark 工作階段,然後重新執行所有儲存格。 |
清理資源
如果你為本文創建了新的湖邊屋,但不再需要:
- 在您的工作區中,以滑鼠右鍵按一下 Lakehouse 名稱。
- 選擇 刪除。
- 確認刪除。
筆記本會留在你的工作區,除非你另外刪除它。