Tareas de clasificación mediante SynapseML

En este artículo se muestra cómo realizar una tarea de clasificación de texto con dos métodos. Un método usa pyspark directamente, y el otro usa la biblioteca synapseml. Ambos métodos producen el mismo rendimiento, pero resaltan cómo SynapseML reduce la complejidad del código en comparación con pyspark.

La tarea consiste en predecir si la reseña de un cliente sobre un libro vendido en Amazon es positiva (calificación > 3) o negativa, basándose en el texto de la reseña. Entrene los alumnos de LogisticRegression con diferentes hiperparámetros y, a continuación, elija el mejor modelo.

Prerrequisitos

  • Cree un notebook.
  • Adjunte su bloc de notas a un lago de datos. En el cuaderno, seleccione Agregar en el panel izquierdo para adjuntar un lago existente o crear uno nuevo.

Note

Todas las bibliotecas que se usan en este artículo (pyspark, synapseml, numpy) están preinstaladas en el entorno de ejecución de Fabric Spark. No es necesario instalar ningún paquete.

Carga y exploración de los datos

En los cuadernos de Fabric, una sesión de Spark ya está disponible como la variable spark. Cargue el conjunto de datos de revisiones de libros de Amazon desde una ubicación pública de Azure Blob Storage:

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

Compruebe que el conjunto de datos se cargó correctamente:

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

Extracción de características y datos de proceso

Los datos reales suelen tener características de varios tipos, por ejemplo, texto, numérico y categórico. Para demostrar cómo trabajar con tipos de características mixtas, agregue dos características numéricas al conjunto de datos: el recuento de palabras de la revisión y la longitud media de la palabra.

Definición de funciones definidas por el usuario (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())

Aplicar UDF con SynapseML UDFTransformer

Use el UDFTransformer de SynapseML para convertir las UDFs en transformadores compatibles con canalizaciones:

from synapse.ml.stages import UDFTransformer

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

Ejecuta la canalización de funciones

Aplique ambos transformadores y cree una columna de etiquetas binarias a partir de la valoración:

from pyspark.ml import Pipeline

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

Compruebe la extracción de características:

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

Clasificación mediante pyspark

Para elegir el mejor clasificador LogisticRegression mediante la pyspark biblioteca, debe realizar explícitamente estos pasos:

  1. Procese las características:
    • Tokenizar la columna de texto.
    • Convierte la columna tokenizada en un vector mediante una función hash.
    • Combine las características numéricas con el vector.
  2. Convierta la columna de etiquetas del tipo booleano al tipo entero.
  3. Entrene varios algoritmos LogisticRegression en el train conjunto de datos con diferentes hiperparámetros.
  4. Calcule el área bajo la curva ROC (AUC) para cada modelo entrenado y seleccione el modelo con la métrica más alta del test conjunto de datos.
  5. Evalúe el mejor modelo en el conjunto validation.

Caracterización y preparación de los datos

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 los datos transformados en características:

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

Entrenamiento y evaluación de 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}")

Compruebe los 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

Los valores exactos de AUC dependen de la división aleatoria. Espere valores entre 0,65 y 0,85.

Clasificación mediante SynapseML

El synapseml enfoque logra el mismo resultado con menos pasos. SynapseML se encarga internamente de la creación de características, lo que reduce el código que hay que escribir:

  1. El TrainClassifier estimador presenta internamente los datos, siempre y cuando las columnas de los trainconjuntos de datos , testy validation representen las características.
  2. El FindBestModel estimador busca el mejor modelo de un grupo de modelos entrenados mediante la evaluación del rendimiento en el test conjunto de datos con la métrica especificada.
  3. El ComputeModelStatistics transformador calcula varias métricas en un conjunto de datos puntuado (en este caso, el validation conjunto de datos) al mismo tiempo.
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)
)

Compruebe los resultados de 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

Los enfoques pyspark y SynapseML deben producir valores AUC similares, ya que entrenan el mismo tipo de modelo con los mismos hiperparámetros en los mismos datos.

Comparación de los dos enfoques

Aspecto pyspark SynapseML
Procesamiento de funcionalidades Manual (de Tokenizer a HashingTF a VectorAssembler) Automático (controlado por TrainClassifier)
Selección de modelos Bucle manual con evaluador Integrado FindBestModel
Cálculo de métricas Métrica única por llamada de evaluación Varias métricas con ComputeModelStatistics
Líneas de código Aproximadamente 30 líneas Aproximadamente 15 líneas
Resultado Misma AUC Misma AUC

Solución de problemas

Cuestión Causa Solución
AnalysisException: Path does not exist La dirección URL de almacenamiento de blobs pública no está disponible temporalmente. Espere unos minutos e inténtelo de nuevo. Comprobación de la conectividad mediante la ejecución de spark.read.parquet("wasbs://publicwasb@mmlspark.blob.core.windows.net/BookReviewsFromAmazon10K.parquet").count()
IllegalArgumentException: Field "features" does not exist Los nombres de columna de características no coinciden entre transformadores Comprobación de los nombres de columna mediante la ejecución data.columns antes del paso VectorAssembler
NameError: name 'LogisticRegression' is not defined Falta la instrucción import Agregar from pyspark.ml.classification import LogisticRegression en la parte superior de la celda
ModuleNotFoundError: No module named 'synapse.ml' El cuaderno no usa el entorno de ejecución de Spark de Fabric Compruebe que el cuaderno usa Fabric Runtime 1.2 o posterior. Seleccione Entorno en la cinta de opciones para comprobarlo.
AUC baja (por debajo de 0,6) Problemas de división de datos o convergencia Compruebe la distribución de etiquetas con data.groupBy("label").count().show(). Espere un conjunto de datos aproximadamente equilibrado.
Py4JJavaError: An error occurred while calling error interno de Java/Spark Compruebe la interfaz de usuario de Spark para obtener registros de errores detallados. Reinicie la sesión de Spark seleccionando Sesión>de detención de sesión y vuelva a ejecutar todas las celdas.

Limpieza de recursos

Si ha creado un nuevo lakehouse para este artículo y ya no lo necesita:

  1. En el área de trabajo, haga clic con el botón derecho en el nombre de lakehouse.
  2. Seleccione Eliminar.
  3. Confirme la eliminación.

El cuaderno permanece en el área de trabajo a menos que lo elimine por separado.