Cómo calcular percentiles y mediana en Apache Spark (ejemplo)
Cómo calcular percentiles y mediana en Apache Spark es una tarea común para resumir distribuciones y detectar outliers. En conjuntos de datos grandes, no es viable traer todos los datos al driver: aquí entran en juego funciones distribuidas como percentile_approx y utilidades como DataFrame.stat.approxQuantile. Esta guía práctica explica por qué elegir cada opción y muestra ejemplos concretos en PySpark, incluyendo consejos de rendimiento y verificación de resultados.
Requisitos previos
- Python 3.7+ y PySpark instalado (versión compatible con tu infraestructura).
- Noções básicas de DataFrame en Apache Spark: crear SparkSession, usar groupBy y agg.
- Entorno local o cluster con memoria suficiente para experimentar. Por ejemplo, para conjuntos de datos de 1–10 millones de filas, se recomienda al menos algunos GB de memoria por executor.
Paso 1: Preparar sesión y datos
Empieza creando una SparkSession y un DataFrame de ejemplo. Usamos un conjunto pequeño para observar el comportamiento; en producción podrías tener millones de filas. Este ejemplo permite probar tanto approxQuantile (ejecutado en el driver) como percentile_approx (distribuido).
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, percentile_approx
spark = SparkSession.builder.appName("QuantilesExample").getOrCreate()
# Dados de exemplo: (group, value)
data = [(1, 10.0),(1, 20.0),(1, 30.0),(2, 5.0),(2, 7.0),(2, 100.0)]
df = spark.createDataFrame(data, ["group","value"])
df.show()
Con estos datos: para el grupo 1 la mediana esperada es 20.0; para el grupo 2 la mediana es 7.0. En total, la distribución tiene un outlier (100.0) que influye en cuartiles y media pero no en la mediana de forma tan severa.
Paso 2: Calcular mediana (percentil 50%) con approxQuantile
La función DataFrame.stat.approxQuantile es sencilla de usar cuando quieres cuantiles de columnas individuales. Ejecuta una operación que devuelve resultados al driver — por eso exige memoria suficiente en el driver cuando aplicas a columnas con muchos valores distintos. El parámetro relativeError controla la tolerancia: 0.01 significa error relativo hasta 1% (valor típico), 0.001 ofrece mayor precisión pero aumenta el coste computacional.
# approxQuantile(col, probabilities, relativeError)
quantiles = df.stat.approxQuantile("value", [0.5], 0.01)
print("Mediana (approxQuantile):", quantiles[0])
# Varios cuantiles: 25%, 50%, 75%
q = df.stat.approxQuantile("value", [0.25, 0.5, 0.75], 0.01)
print("Cuartiles:", q)
En conjuntos de datos grandes (por ejemplo, 10M+ filas) evita usar approxQuantile para muchas columnas a la vez sin garantizar memoria en el driver; usa muestreo o percentile_approx para agregaciones.
Paso 3: Usar percentile_approx para agregaciones y por grupo
percentile_approx está implementado como una función de agregación distribuida — ideal para groupBy. Acepta un único percentil o una lista y devuelve un array cuando pides múltiples percentiles. El tercer argumento (accuracy) controla la granularidad del resumen interno: valores típicos van de 100 a 10000; cuanto mayor, mayor precisión y coste.
from pyspark.sql.functions import percentile_approx
# Mediana global con percentile_approx
df.agg(percentile_approx(col("value"), 0.5).alias("median")).show()
# Cuartiles globales (retorna array)
df.agg(percentile_approx("value", [0.25, 0.5, 0.75]).alias("quartis")).show()
# Cuartiles por grupo
quartis_por_grupo = df.groupBy("group").agg(
percentile_approx("value", [0.25, 0.5, 0.75]).alias("quartis")
)
quartis_por_grupo.show()
# Separar los cuartiles en columnas legibles
from pyspark.sql.functions import col
quartis_por_grupo.select(
col("group"),
col("quartis").getItem(0).alias("q1"),
col("quartis").getItem(1).alias("q2"),
col("quartis").getItem(2).alias("q3")
).show()
Ejemplo práctico: en un conjunto con 6 valores, pedir [0.25, 0.5, 0.75] devuelve un array con tres entradas. Para grupos muy desbalanceados (por ejemplo, un grupo con 10M filas y otro con 10 filas) puedes observar diferencias en la precisión entre grupos; considera ajustar accuracy caso por caso.
Paso 4: Manejar valores nulos, tipos y precisión
Antes de calcular percentiles, trata nulos y garantiza tipos numéricos. Convierte a DoubleType si hay enteros y decimales mezclados. Para aumentar precisión en percentile_approx, especifica un valor de accuracy mayor (ej.: 1000 o 10000), pero esto aumenta el uso de CPU y memoria.
# Remover nulos
df_clean = df.filter(col("value").isNotNull())
# Convertir tipo (ejemplo)
df_clean = df_clean.withColumn("value", col("value").cast("double"))
# Ejemplo con mayor precisión (accuracy = 10000)
df_clean.agg(percentile_approx("value", 0.5, 10000).alias("median_preciso")).show()
Si trabajas con datos financieros, considera usar DecimalType para preservar precisión; sin embargo, algunas funciones pueden exigir cast a double.
Verificar el resultado
Valida resultados comparando métodos: approxQuantile (lista en el driver) versus percentile_approx (DataFrame/array). Para conjuntos pequeños, calcula la mediana exacta con Python local para confirmar. En grandes conjuntos de datos, empieza con relativeError ~0.01 o accuracy ~1000 y solo aumenta si las pruebas indican sesgo significativo.
# Mediana exacta en datos pequeños (solo para verificación)
vals = [r[0] for r in df.select("value").rdd.flatMap(lambda x: x).collect()]
vals.sort()
from statistics import median
print("Mediana exacta (muestra):", median(vals))
Conclusión
Has calculado percentiles y mediana en Apache Spark usando approxQuantile para columnas únicas y percentile_approx para agregaciones y groupBy. En resumen: usa approxQuantile cuando puedes traer un resumen al driver y necesitas cuantiles rápidos para una columna; usa percentile_approx para operaciones distribuidas y por grupo. Ajusta relativeError y accuracy según el compromiso entre precisión y coste. Próximos pasos: experimentar con percentiles múltiples en datos reales (ej.: 1M–100M filas), medir tiempo y memoria y documentar la precisión que consideres aceptable para tu caso.