Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
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
Obtenga una suscripción a Microsoft Fabric. O bien, regístrese para obtener una evaluación gratuita de Microsoft Fabric.
Inicie sesión en Microsoft Fabric.
Cambie a Fabric mediante el conmutador de experiencia en el lado inferior izquierdo de la página principal.
- 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:
- 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.
- Convierta la columna de etiquetas del tipo booleano al tipo entero.
- Entrene varios algoritmos LogisticRegression en el
trainconjunto de datos con diferentes hiperparámetros. - Calcule el área bajo la curva ROC (AUC) para cada modelo entrenado y seleccione el modelo con la métrica más alta del
testconjunto de datos. - 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:
- El
TrainClassifierestimador presenta internamente los datos, siempre y cuando las columnas de lostrainconjuntos de datos ,testyvalidationrepresenten las características. - El
FindBestModelestimador busca el mejor modelo de un grupo de modelos entrenados mediante la evaluación del rendimiento en eltestconjunto de datos con la métrica especificada. - El
ComputeModelStatisticstransformador calcula varias métricas en un conjunto de datos puntuado (en este caso, elvalidationconjunto 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:
- En el área de trabajo, haga clic con el botón derecho en el nombre de lakehouse.
- Seleccione Eliminar.
- Confirme la eliminación.
El cuaderno permanece en el área de trabajo a menos que lo elimine por separado.