Como calcular percentis e mediana em Apache Spark (exemplo)
Como calcular percentis e mediana em Apache Spark é uma tarefa comum para sumarizar distribuições e detetar outliers. Em conjuntos de dados grandes, não é viável trazer todos os dados para o driver: é aqui que entram funções distribuídas como percentile_approx e utilitários como DataFrame.stat.approxQuantile. Este guia prático explica o porquê de cada opção e mostra exemplos concretos em PySpark, incluindo dicas de desempenho e verificação de resultados.
Pré-requisitos
- Python 3.7+ e PySpark instalado (versão compatível com a tua infra).
- Noções básicas de DataFrame em Apache Spark: criar SparkSession, usar groupBy e agg.
- Ambiente local ou cluster com memória suficiente para experimentar. Por exemplo, para conjuntos de dados de 1–10 milhões de linhas, recomenda-se pelo menos alguns GB de memória por executor.
Passo 1: Preparar sessão e dados
Começa por criar uma SparkSession e um DataFrame de exemplo. Usamos um pequeno conjunto para ver o comportamento; em produção poderás ter milhões de linhas. Este exemplo permite testar tanto approxQuantile (executado no driver) como percentile_approx (distribuído).
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()
Com estes dados: para o grupo 1 a mediana esperada é 20.0; para o grupo 2 a mediana é 7.0. No total, a distribuição tem um outlier (100.0) que influencia quartis e média mas não a mediana de forma tão severa.
Passo 2: Calcular mediana (percentil 50%) com approxQuantile
A função DataFrame.stat.approxQuantile é simples de usar quando queres quantis de colunas individuais. Ela executa uma operação que devolve resultados ao driver — por isso exige memória suficiente no driver quando aplicas a colunas com muitos valores distintos. O parâmetro relativeError controla a tolerância: 0.01 significa erro relativo até 1% (valor típico), 0.001 oferece maior precisão mas aumenta o custo computacional.
# approxQuantile(col, probabilities, relativeError)
quantiles = df.stat.approxQuantile("value", [0.5], 0.01)
print("Mediana (approxQuantile):", quantiles[0])
# Vários quantis: 25%, 50%, 75%
q = df.stat.approxQuantile("value", [0.25, 0.5, 0.75], 0.01)
print("Quartis:", q)
Em conjuntos de dados grandes (por exemplo, 10M+ linhas) evita usar approxQuantile para muitas colunas em simultâneo sem garantir memória no driver; usa amostragem ou percentile_approx para agregações.
Passo 3: Usar percentile_approx para agregações e por grupo
percentile_approx é implementado como uma função de agregação distribuída — ideal para groupBy. Aceita um único percentil ou uma lista e devolve um array quando pedes múltiplos percentis. O terceiro argumento (accuracy) controla a granularidade do resumo interno: valores típicos vão de 100 a 10000; quanto maior, maior precisão e custo.
from pyspark.sql.functions import percentile_approx
# Mediana global com percentile_approx
df.agg(percentile_approx(col("value"), 0.5).alias("median")).show()
# Quartis globais (retorna array)
df.agg(percentile_approx("value", [0.25, 0.5, 0.75]).alias("quartis")).show()
# Quartis 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 os quartis em colunas legíveis
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()
Exemplo prático: num conjunto com 6 valores, pedir [0.25, 0.5, 0.75] devolve um array com três entradas. Para grupos muito desbalanceados (por exemplo, um grupo com 10M linhas e outro com 10 linhas) podes ver diferenças na precisão entre grupos; considera ajustar accuracy caso a caso.
Passo 4: Lidar com valores nulos, tipos e precisão
Antes de calcular percentis, trata nulos e garante tipos numéricos. Converte para DoubleType se houver inteiros e decimais mistos. Para aumentar precisão em percentile_approx, especifica um valor de accuracy maior (ex.: 1000 ou 10000), mas isso aumenta o uso de CPU e de memória.
# Remover nulos
df_clean = df.filter(col("value").isNotNull())
# Converter tipo (exemplo)
df_clean = df_clean.withColumn("value", col("value").cast("double"))
# Exemplo com maior precisão (accuracy = 10000)
df_clean.agg(percentile_approx("value", 0.5, 10000).alias("median_preciso")).show()
Se trabalhares com dados financeiros, considera usar DecimalType para preservar precisão; no entanto, algumas funções podem exigir cast para double.
Verificar o resultado
Valida resultados comparando métodos: approxQuantile (lista no driver) versus percentile_approx (DataFrame/array). Para conjuntos pequenos, calcula a mediana exata com Python local para confirmar. Em grandes conjuntos de dados, começa com relativeError ~0.01 ou accuracy ~1000 e só aumenta se os testes indicarem viés significativo.
# Mediana exata em dados pequenos (apenas para verificação)
vals = [r[0] for r in df.select("value").rdd.flatMap(lambda x: x).collect()]
vals.sort()
from statistics import median
print("Mediana exata (amostra):", median(vals))
Conclusão
Calculaste percentis e mediana em Apache Spark usando approxQuantile para colunas únicas e percentile_approx para agregações e groupBy. Em resumo: usa approxQuantile quando podes trazer um resumo para o driver e precisas de quantis rápidos para uma coluna; usa percentile_approx para operações distribuídas e por grupo. Ajusta relativeError e accuracy conforme o compromisso entre precisão e custo. Próximos passos: experimentar com percentis múltiplos em dados reais (ex.: 1M–100M linhas), medir tempo e memória e documentar a precisão que consideras aceitável para o teu caso.