SynapseML を使用した分類タスク

この記事では、2 つの方法でテキスト分類タスクを実行する方法について説明します。 1 つの方法ではプレーン pysparkを使用し、もう 1 つの方法では synapseml ライブラリを使用します。 どちらの方法でも同じパフォーマンスが得られますが、SynapseML が pysparkと比較してコードの複雑さを軽減する方法を強調します。

このタスクでは、Amazon で販売された書籍の顧客レビューが、レビュー テキストに基づいて良好 (評価 > 3) か悪いかを予測します。 さまざまなハイパーパラメーターを使用して LogisticRegression 学習者をトレーニングし、最適なモデルを選択します。

[前提条件]

  • ノートブックを作成 します
  • ノートブックをレイクハウスにアタッチします。 ノートブックで、左側のウィンドウで [追加 ] を選択して既存のレイクハウスをアタッチするか、新しいレイクハウスを作成します。

Note

この記事で使用されるすべてのライブラリ (pysparksynapsemlnumpy) は、Fabric Spark ランタイムにプレインストールされています。 パッケージをインストールする必要はありません。

データを読み込んで探索する

Fabric ノートブックでは、Spark セッションは既に spark 変数として使用できます。 パブリック Azure Blob Storageの場所から Amazon 書籍レビュー データセットを読み込みます。

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

特徴の抽出とデータの処理

多くの場合、実際のデータには、テキスト、数値、カテゴリなどの複数の型の特徴があります。 混在する特徴の種類の操作を示すには、レビューの 単語数平均単語の長さという 2 つの数値特徴をデータセットに追加します。

ユーザー定義関数 (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 を使用して分類する

pyspark ライブラリを使用して最適な LogisticRegression 分類子を選択するには、次の手順を明示的に実行する必要があります。

  1. 機能を処理します。
    • テキスト列をトークン化します。
    • ハッシュを使用して、トークン化された列をベクターにハッシュします。
    • 数値特徴をベクトルとマージします。
  2. ラベル列をブール型から整数型にキャストします。
  3. 異なるハイパーパラメーターを使用して、 train データセットで複数の LogisticRegression アルゴリズムをトレーニングします。
  4. トレーニング済みモデルごとに ROC 曲線下面積 (AUC) を計算し、 test データセットで最も高いメトリックを持つモデルを選択します。
  5. 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 では、内部的に特徴付けが処理されるため、記述する必要があるコードが減ります。

  1. TrainClassifierエスティメーターは、traintest、およびvalidationデータセット内の列が特徴を表す限り、データを内部的に特徴付けします。
  2. FindBestModel推定器は、指定されたメトリックを使用してtest データセットのパフォーマンスを評価することで、トレーニング済みモデルのプールから最適なモデルを検索します。
  3. 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 値が生成されます。

2 つの方法を比較する

特徴 pyspark SynapseML
特徴処理 手動(Tokenizer → HashingTF → VectorAssembler) 自動 ( TrainClassifierによって処理されます)
モデルの選択 エバリュエーターを使用した手動ループ 組み込み FindBestModel
メトリックの計算 評価呼び出しごとに 1 つのメトリック 複数のメトリック ComputeModelStatistics
コード行 約 30 行 約 15 行
結果 同じ AUC 同じ AUC

Troubleshooting

問題 原因 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 import ステートメントがありません セルの上部に from pyspark.ml.classification import LogisticRegression を追加する
ModuleNotFoundError: No module named 'synapse.ml' ノートブックで Fabric Spark ランタイムを使用していません ノートブックでランタイム 1.2 以降Fabric使用されていることを確認します。 確認するには、リボンの [ 環境 ] を選択します。
低 AUC (0.6 未満) データ分割の問題または収束問題 data.groupBy("label").count().show()を使用してラベルの分布を確認します。 ほぼバランスの取れたデータセットが必要です。
Py4JJavaError: An error occurred while calling Java/Spark 内部エラー Spark UI で詳細なエラー ログを確認します。 [セッション>ストップ セッション] を選択して Spark セッションを再起動し、すべてのセルを再実行します。

リソースをクリーンアップする

この記事用に新しいレイクハウスを作成し、不要になった場合:

  1. ワークスペースで、lakehouse 名を右クリックします。
  2. を選択して、を削除します。
  3. 削除を確認します。

ノートブックは、個別に削除しない限り、ワークスペースに残ります。