Tarefas de classificação usando SynapseML

Este artigo mostra como realizar uma tarefa de classificação de texto com dois métodos. Um método usa o simples pyspark, e o outro usa a synapseml biblioteca. Ambos os métodos produzem o mesmo desempenho, mas destacam como o SynapseML reduz a complexidade do código em comparação com pyspark.

A tarefa prevê se uma avaliação de um cliente de um livro vendido na Amazon é boa (classificação > 3) ou má, com base no texto da crítica. Treinas aprendizes LogisticRegression com diferentes hiperparâmetros, e depois escolhes o melhor modelo.

Pré-requisitos

  • Crie um bloco de anotações.
  • Ligue o seu bloco de notas a uma casa no lago. No caderno, selecione Adicionar no painel esquerdo para anexar uma casa de lago existente ou criar uma nova.

Note

Todas as bibliotecas usadas neste artigo (pyspark, synapseml, numpy) estão pré-instaladas no runtime Fabric Spark. Não precisas de instalar nenhum pacote.

Carregue e explore os dados

Nos cadernos do Fabric, uma sessão Spark já está disponível na variável spark. Carregue o conjunto de dados de críticas de livros da Amazon a partir de uma localização pública no Armazenamento de Blobs do Azure:

rawData = spark.read.parquet(
    "wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet"
)
rawData.show(5)

Verifique se o conjunto de dados foi carregado corretamente:

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

Extrair características e processar dados

Os dados reais frequentemente apresentam características de vários tipos, por exemplo, texto, numérico e categórico. Para demonstrar o trabalho com tipos de características mistos, adicione duas características numéricas ao conjunto de dados: a contagem de palavras da revisão e o comprimento médio das palavras.

Definir funções definidas pelo utilizador (UDFs)

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

Aplicar UDFs com SynapseML UDFTransformer

Utilize o UDFTransformer do SynapseML para envolver os UDFs em transformadores compatíveis com pipelines:

from synapse.ml.stages import UDFTransformer

wordLengthTransformer = UDFTransformer(
    inputCol="text", outputCol="wordLength", udf=wordLengthUDF
)
wordCountTransformer = UDFTransformer(
    inputCol="text", outputCol="wordCount", udf=wordCountUDF
)

Executar o pipeline de funcionalidades

Aplique ambos os transformadores e crie uma coluna de etiqueta binária a partir da classificação:

from pyspark.ml import Pipeline

data = (
    Pipeline(stages=[wordLengthTransformer, wordCountTransformer])
    .fit(rawData)
    .transform(rawData)
    .withColumn("label", rawData["rating"] > 3)
    .drop("rating")
)

Verifique a extração da funcionalidade:

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

Classificar usando pyspark

Para escolher o melhor classificador LogisticRegression usando a pyspark biblioteca, deve executar explicitamente estes passos:

  1. Processe as funcionalidades:
    • Divida a coluna de texto em tokens.
    • Transforme a coluna tokenizada num vetor usando hashing.
    • Juntar as características numéricas com o vetor.
  2. Converta a coluna de rótulo do tipo booleano para o tipo inteiro.
  3. Treinar vários algoritmos LogisticRegression no conjunto de dados train com diferentes hiperparâmetros.
  4. Calcule a Área Sob a Curva ROC (AUC) para cada modelo treinado e selecione o modelo com a métrica mais alta no test conjunto de dados.
  5. Avalie o melhor modelo no conjunto validation.

Personaliza e prepara os dados

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

Verifique os dados característicos:

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

Treinar e avaliar modelos

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

Verifique os resultados:

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

Os valores exatos do AUC dependem da divisão aleatória. Espere valores entre 0,65 e 0,85.

Classificar usando SynapseML

A synapseml abordagem alcança o mesmo resultado com menos passos. O SynapseML trata internamente da extração de características, o que reduz a quantidade de código que precisa de escrever:

  1. O TrainClassifier estimador apresenta internamente os dados, desde que as colunas nos trainconjuntos de dados , test, e validation representem as características.
  2. O FindBestModel estimador encontra o melhor modelo a partir de um conjunto de modelos treinados, avaliando o desempenho no test conjunto de dados com a métrica especificada.
  3. O ComputeModelStatistics transformador calcula múltiplas métricas num conjunto de dados pontuado (neste caso, o validation conjunto de dados) ao mesmo tempo.
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)
)

Verifique os resultados do 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

As abordagens pyspark e SynapseML devem produzir valores AUC semelhantes, uma vez que treinam o mesmo tipo de modelo com os mesmos hiperparâmetros nos mesmos dados.

Compare as duas abordagens.

Aspect Pyspark SynapseML
Processamento de funcionalidades Manual (Tokenizer para HashingTF para VectorAssembler) Automático (tratado por TrainClassifier)
Seleção de modelos Ciclo manual com avaliador Integrado FindBestModel
Cálculo de métricas Métrica única por chamada de avaliação Múltiplas métricas com ComputeModelStatistics
Linhas de código Cerca de 30 linhas Cerca de 15 linhas
Result Mesma AUC Mesma AUC

Troubleshooting

Problema Motivo Resolução
AnalysisException: Path does not exist A URL pública de armazenamento de blob está temporariamente indisponível Aguarde alguns minutos e tente novamente. Verifique a conectividade executando spark.read.parquet("wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet").count()
IllegalArgumentException: Field "features" does not exist Os nomes das colunas de atributos não coincidem entre os transformadores Verifique os nomes das colunas executando data.columns antes do passo VectorAssembler
NameError: name 'LogisticRegression' is not defined Extrato de importação em falta Adicione from pyspark.ml.classification import LogisticRegression no topo da célula
ModuleNotFoundError: No module named 'synapse.ml' O notebook não está a usar o runtime do Fabric Spark Verifique se o portátil usa o Fabric Runtime 1.2 ou posterior. Selecione Ambiente na fita para verificar.
AUC baixo (abaixo de 0,6) Questão de divisão de dados ou problemas de convergência Verifique a distribuição do rótulo com data.groupBy("label").count().show(). Espere um conjunto de dados aproximadamente equilibrado.
Py4JJavaError: An error occurred while calling Erro interno Java/Spark Consulta a interface do Spark para registos de erros detalhados. Reinicie a sessão do Spark ao selecionar Sessão>Parar sessão e, em seguida, execute novamente todas as células.

Limpeza de recursos

Se criou uma nova casa no lago para este artigo e já não precisa dela:

  1. No seu espaço de trabalho, clique com o botão direito no nome da casa do lago.
  2. Selecione Eliminar.
  3. Confirme a exclusão.

O caderno permanece no seu espaço de trabalho, a menos que o apague separadamente.