(+351) 21 24 10006  ·  info@bconcepts.pt
Carnaxide, Lisboa

Cómo calcular y aplicar Z-score por grupo en Apache Spark: paso a paso

João Barros 19 de September de 2026 5 min de lectura

Normalizar datos por grupo ayuda a comparar valores dentro de categorías distintas y a detectar anomalías relativas. Este tutorial muestra cómo calcular y aplicar el Z-score por grupo en Apache Spark (PySpark) para transformar una variable numérica según la media y la desviación estándar de su grupo.

Requisitos previos

  • PySpark instalado y configurado y una SparkSession funcional (por ejemplo, Spark 3.x).
  • Conocimientos básicos de DataFrame, SQL y funciones de agregación en Spark (groupBy, agg, join, Window).
  • Editor o Notebook (Jupyter, VS Code, Databricks) para ejecutar código Python e inspeccionar resultados.
  • Preferible: entender las diferencias entre stddev_pop y stddev_samp — esto influye en la interpretación de la desviación estándar.

Paso 1: Entender el objetivo y los motivos

El Z-score transforma cada valor x en (x - mean)/std. Hacer esto por grupo (por ejemplo, por producto, tienda o región) permite comparar valores relativos dentro de cada grupo independientemente de la escala absoluta. Por ejemplo, ventas de 110 unidades pueden ser normales en una región donde la media es 100 con std ≈ 5 (z ≈ 2), pero serían un outlier en una región con media 105 y std ≈ 1 (z ≈ 5).

Usos prácticos: detectar anomalías por producto, crear características escaladas para modelos de machine learning, o visualizar desviaciones normalizadas entre categorías. Reduce sesgo cuando los grupos tienen escalas muy diferentes (p. ej.: precios en distintos mercados) y facilita reglas operativas (por ejemplo, investigar puntos con |z| > 3).

Paso 2: Crear un ejemplo mínimo de datos

Crear un DataFrame pequeño permite validar la lógica antes de aplicar a grandes volúmenes. Aquí usamos dos grupos con cuatro registros cada uno — suficiente para calcular media y stdpop. Valores reales tendrán miles a millones de filas, pero la lógica es equivalente.

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

spark = SparkSession.builder.appName("zscore-group-example").getOrCreate()

data = [
    ("A", 10.0), ("A", 12.0), ("A", 8.0), ("A", 25.0),
    ("B", 100.0), ("B", 110.0), ("B", 95.0), ("B", 104.0),
]

df = spark.createDataFrame(data, ["group", "value"]) 
df.show()

Paso 3: Calcular media y desviación estándar por grupo

Agrupa por grupo para obtener mean y std. Usa funciones built-in para garantizar rendimiento en clusters. Elige stddev_pop si consideras la población completa del grupo, o stddev_samp si tienes una muestra. Para nuestro ejemplo, usando stddev_pop, los resultados son aproximados:

  • Grupo A: mean ≈ 13.75, std ≈ 6.645 (calculado sobre 4 observaciones).
  • Grupo B: mean ≈ 102.25, std ≈ 5.494.

Estos números permiten verificar después los Z-scores manualmente (p. ej.: para A, valor 25 → z ≈ (25-13.75)/6.645 ≈ 1.69).

from pyspark.sql.functions import mean, stddev_pop

stats = df.groupBy("group").agg(
    mean("value").alias("mean_value"),
    stddev_pop("value").alias("std_value")
)

stats.show()

Paso 4: Unir estadísticas al DataFrame original

Haz un join entre el DataFrame original y las estadísticas por grupo. Si las estadísticas (stats) son pequeñas comparadas con el dataset, considera un broadcast join para eficiencia: broadcast(stats). Para conjuntos grandes, el join por clave es el estándar. Alternativa: calcular media y std con Window functions (over partitionBy) para evitar un join explícito — útil en pipelines donde prefieras evitar pasos de shuffle adicionales.

df_with_stats = df.join(stats, on="group", how="left")
df_with_stats.show()

Paso 5: Calcular el Z-score y tratar casos de desviación estándar cero

Calcula (value - mean_value) / std_value. Grupos con un único registro o sin variabilidad tendrán std = 0; evita división por cero reemplazando por un comportamiento definido: poner z_score = 0, NaN, o usar un pequeño epsilon. Otra aproximación es exigir count >= 2 para calcular Z-score y señalar los restantes como "insuficientes".

from pyspark.sql.functions import when, lit

epsilon = 1e-9

result = df_with_stats.withColumn(
    "z_score",
    when(col("std_value").isNull() | (col("std_value") == 0), lit(0.0))
    .otherwise((col("value") - col("mean_value")) / (col("std_value") + lit(epsilon)))
)

result.select("group", "value", "mean_value", "std_value", "z_score").show()

Paso 6: Usos prácticos — detectar anomalías y filtrar

Un umbral común para anomalías es |z_score| > 3: en una distribución normal esto corresponde a ~0.27% de los puntos. Si prefieres mayor sensibilidad, usa 2.5 (~1.24%) o 2 (~5%). El umbral debe reflejar el coste de falsos positivos vs falsos negativos en tu contexto. Por ejemplo, en una operación con 1M de registros por día, aplicar |z|>3 puede devolver ~2700 registros sospechosos por día (suposición de normalidad), un volumen manejable para revisión manual.

threshold = 3.0
anomalies = result.filter(abs(col("z_score")) > threshold)
anomalies.show()

Aparte de la detección, guarda la columna z_score como característica para modelos (normaliza por grupo antes de entrenar), o calcula percentiles por grupo para reglas de negocio. Prueba siempre con grupos de tamaños variados y valida la tasa de alarmas con datos etiquetados, si existen.

Verificar el resultado

Confirma que las medias y desviaciones estándar por grupo tienen sentido usando stats.show(). Verifica manualmente algunos cálculos (por ejemplo, los ejemplos numéricos dados) y prueba grupos con un único registro para confirmar el tratamiento de la división por cero. Si los resultados son extraños, revisa agregaciones y valores nulos.

Conclusión

Normalizar por grupo con Z-score en Apache Spark es una técnica simple pero potente para detección de anomalías y preparación de características. Para producción, incorpora este cálculo en un pipeline ETL, valida umbrales con datos históricos y elige stddev_pop vs stddev_samp según la definición estadística necesaria. Paso siguiente: aplicar por ventanas temporales (p. ej.: Z-score por grupo y por mes) o integrar en una etapa de validación automática.