Cómo detectar y eliminar outliers en Apache Spark: paso a paso
Este tutorial muestra cómo detectar y eliminar outliers en Apache Spark para mejorar la calidad de los datos y el rendimiento de modelos predictivos. Vamos a explicar por qué funcionan los enfoques IQR y Z-score, cuándo cada uno es más apropiado y cómo aplicar esas técnicas de forma escalable en PySpark para conjuntos de datos con cientos de miles o millones de filas.
Prerequisitos
- Instalación de Spark y PySpark disponible (local o clúster). Se recomienda Spark 3.x para mejores funciones y rendimiento.
- Archivo CSV o Parquet con columnas numéricas para análisis; para conjuntos de datos grandes (p. ej.: 1M–50M filas) usa Parquet para lectura más rápida y menor I/O.
- Conocimientos básicos de DataFrame y funciones de PySpark. Conocimientos básicos sobre medias, desviación estándar y percentiles ayudan a interpretar resultados.
- Configuración de recursos: si tienes 10 ejecutores con 4 cores cada uno, la aproximación de cuantiles (approxQuantile) es normalmente rápida; si trabajáis con 100M+ filas ajusta el parámetro relativeError.
Paso 1: Cargar los datos en PySpark
Comenzamos leyendo los datos a un DataFrame. Usa Parquet siempre que sea posible (más rápido y mantiene tipos). Si cargas CSV, especifica schema para evitar inferencia costosa. Verifica y trata los nulls antes de calcular estadísticas — por ejemplo, en muchos escenarios 0.1–2% de las filas pueden tener nulls en columnas numéricas y conviene decidir si rellenar o descartar.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col
spark = SparkSession.builder.appName("outliers-example").getOrCreate()
df = spark.read.csv("/caminho/dados.csv", header=True, inferSchema=True)
# Selecionar colunas numéricas de interesse
numeric_cols = ["valor1", "valor2"]
df = df.select([col(c) for c in numeric_cols])
df.printSchema()
Ejemplo práctico: en un conjunto de datos con 5M de filas, df.count() puede tardar desde segundos hasta minutos dependiendo del clúster; evita count() innecesarios en producción.
Paso 2: Detectar outliers con IQR (Interquartile Range)
IQR es robusto frente a valores extremos y no asume normalidad. Calcula Q1 y Q3 por columna, define límites [Q1 - 1.5*IQR, Q3 + 1.5*IQR] (1.5 es la regla clásica) y marca filas con valores fuera de esos límites. Para conjuntos de datos muy grandes usa approxQuantile con relativeError típicamente entre 0.01 y 0.001. Por ejemplo, con 10M de filas y relativeError=0.01 normalmente se obtienen cuantiles dentro del 1%.
from pyspark.sql.functions import expr
# Calcular cuantiles aproximados (más rápido para datos grandes)
quantiles = df.approxQuantile(numeric_cols, [0.25, 0.75], 0.01)
# quantiles es lista de pares [q1, q3] por columna
bounds = {}
for i, col_name in enumerate(numeric_cols):
q1, q3 = quantiles[i]
iqr = q3 - q1
lower = q1 - 1.5 * iqr
upper = q3 + 1.5 * iqr
bounds[col_name] = (lower, upper)
# Crear expresión para filtrar outliers por columna
outlier_cond = " OR ".join([f"{c} < {bounds[c][0]} OR {c} > {bounds[c][1]}" for c in numeric_cols])
outliers_iqr = df.filter(expr(outlier_cond))
non_outliers_iqr = df.filter(~expr(outlier_cond))
En la práctica, en muchos casos solo el 0.1%–5% de los registros son marcados como outliers con IQR; ese valor depende del dominio (p. ej.: los sensores tienen más ruido que las transacciones financieras).
Paso 3: Detectar outliers con Z-score (supuesto de normalidad)
Z-score es útil cuando la distribución es aproximadamente normal. Calcula media y desviación estándar, transforma valores en Z y marca puntos con |Z| > threshold. Un threshold común es 3 (aproximadamente 0.3% de los puntos en una distribución normal). Atención: si la desviación estándar es cero, Z-score no es aplicable — normalmente implica columna constante o datos incorrectos.
from pyspark.sql.functions import mean, stddev
stats = df.agg(*[mean(c).alias(c+"_mean") for c in numeric_cols], *[stddev(c).alias(c+"_std") for c in numeric_cols]).collect()[0]
z_exprs = []
threshold = 3.0
for c in numeric_cols:
mu = stats[c+"_mean"]
sigma = stats[c+"_std"] if stats[c+"_std"] is not None else 0.0
if sigma == 0.0:
# No es posible calcular Z-score si la desviación es cero; evitar división por cero
z_exprs.append(f"false")
else:
z_exprs.append(f"abs(({c} - {mu}) / {sigma}) > {threshold}")
outlier_cond_z = " OR ".join(z_exprs)
outliers_z = df.filter(expr(outlier_cond_z))
non_outliers_z = df.filter(~expr(outlier_cond_z))
Comparar IQR vs Z-score: en distribuciones con colas largas el IQR tiende a marcar menos falsos positivos; con distribuciones casi gaussianas Z-score es sensible e interpretable (|Z|>3).
Paso 4: Elegir estrategia y eliminar outliers
Decide si quieres eliminar todos los outliers identificados por cualquier método o usar solo uno. Estrategias comunes: (a) eliminar por la unión de métodos (más conservador), (b) eliminar por la intersección (solo si ambos coinciden), (c) aplicar clipping o winsorization en lugar de borrar filas. Por ejemplo, si eliminas 2% de 1M de filas te quedas con 980k; esa pérdida puede ser aceptable, pero si eliminas 20% conviene revisar los thresholds.
# Remover outliers IQR (ejemplo)
df_clean = non_outliers_iqr
# Alternativa: remover unión de outliers IQR y Z-score
# outliers_union = outliers_iqr.union(outliers_z).dropDuplicates()
# df_clean = df.join(outliers_union, on=numeric_cols, how='left_anti')
# Grabar resultado
df_clean.write.mode("overwrite").parquet("/caminho/dados_clean.parquet")
Si prefieres no perder registros, puedes aplicar clipping: por ejemplo truncar valores por debajo del lower a lower y por encima del upper a upper, lo que mantiene el número de filas pero reduce el impacto de los extremos.
Verificar el resultado
Confirma que los outliers fueron eliminados mostrando conteos y estadísticas antes/después. Observa percentiles clave (1%, 50%, 99%) para confirmar reducción de extremos. Es importante validar también el impacto en el modelo: entrena el modelo antes y después y compara métricas (RMSE, AUC, etc.) — a veces la eliminación mejora el rendimiento, en otras ocasiones puedes perder señal útil.
print("Total original:", df.count())
print("Total limpio:", df_clean.count())
# Estadísticas antes y después
print("Estadísticas originales:")
df.describe().show()
print("Estadísticas limpias:")
df_clean.describe().show()
# Ver cuantiles para confirmar eliminación de extremos
print("Cuantiles limpios:", df_clean.approxQuantile(numeric_cols, [0.01, 0.5, 0.99], 0.01))
Conclusión
Has borrado o tratado outliers usando IQR y Z-score en Apache Spark. Estas técnicas ayudan a mejorar análisis y modelos cuando se aplican con cuidado. Próximos pasos prácticos: experimentar con thresholds (p. ej.: 1.5→3 para IQR, 3→4 para Z), probar clipping vs eliminación, y validar el impacto en tu modelo con datos de validación. Consejo: registra cuántos registros se eliminan y conserva una copia de los outliers para investigación posterior.