Como calcular e aplicar Z-score por grupo em Apache Spark: passo a passo
Normalizar dados por grupo ajuda a comparar valores dentro de categorias distintas e a detetar anomalias relativas. Este tutorial mostra como calcular e aplicar o Z-score por grupo em Apache Spark (PySpark) para transformar uma variável numérica segundo a média e desvio-padrão do seu grupo.
Pré-requisitos
- PySpark instalado e configurado e uma SparkSession funcional (por exemplo, Spark 3.x).
- Conhecimentos básicos de DataFrame, SQL e funções de agregação em Spark (groupBy, agg, join, Window).
- Editor ou Notebook (Jupyter, VS Code, Databricks) para executar código Python e inspecionar resultados.
- Preferível: perceber as diferenças entre stddev_pop e stddev_samp — isto influencia a interpretação do desvio-padrão.
Passo 1: Entender o objetivo e os motivos
O Z-score transforma cada valor x em (x - mean)/std. Fazer isto por grupo (por exemplo, por produto, loja ou região) permite comparar valores relativos dentro de cada grupo independentemente da escala absoluta. Por exemplo, vendas de 110 unidades podem ser normais numa região onde a média é 100 com std ≈ 5 (z ≈ 2), mas seriam um outlier numa região com média 105 e std ≈ 1 (z ≈ 5).
Usos práticos: detetar anomalias por produto, criar características escaladas para modelos de machine learning, ou visualizar desvios normalizados entre categorias. Reduz viés quando grupos têm escalas muito diferentes (ex.: preços em diferentes mercados) e facilita regras operacionais (por exemplo, investigar pontos com |z| > 3).
Passo 2: Criar um exemplo mínimo de dados
Criar um DataFrame pequeno permite validar a lógica antes de aplicar a grandes volumes. Aqui usamos dois grupos com quatro registos cada — suficiente para calcular média e stdpop. Valores reais terão milhares a milhões de linhas, mas a lógica é 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()
Passo 3: Calcular média e desvio-padrão por grupo
Agrega por grupo para obter mean e std. Usa funções built-in para garantir performance em clusters. Escolhe stddev_pop se consideras a população completa do grupo, ou stddev_samp se tens uma amostra. Para o nosso exemplo, usando stddev_pop, os resultados são aproximados:
- Grupo A: mean ≈ 13.75, std ≈ 6.645 (calculado sobre 4 observações).
- Grupo B: mean ≈ 102.25, std ≈ 5.494.
Estes números permitem verificar depois os Z-scores manualmente (ex.: 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()
Passo 4: Unir estatísticas ao DataFrame original
Faz um join entre o DataFrame original e as estatísticas por grupo. Se as estatísticas (stats) forem pequenas comparadas ao dataset, considera um broadcast join para eficiência: broadcast(stats). Para grandes conjuntos, o join por chave é o padrão. Alternativa: calcular média e std com Window functions (over partitionBy) para evitar um join explícito — útil em pipelines onde preferes evitar etapas de shuffle adicionais.
df_with_stats = df.join(stats, on="group", how="left")
df_with_stats.show()
Passo 5: Calcular o Z-score e tratar casos de desvio-padrão zero
Calcula (value - mean_value) / std_value. Grupos com um único registo ou sem variabilidade terão std = 0; evita divisão por zero substituindo por um comportamento definido: colocar z_score = 0, NaN, ou usar um pequeno epsilon. Outra abordagem é exigir count >= 2 para calcular Z-score e sinalizar os 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()
Passo 6: Usos práticos — detetar anomalias e filtrar
Um limiar comum para anomalias é |z_score| > 3: numa distribuição normal isso corresponde a ~0.27% dos pontos. Se preferires maior sensibilidade, usa 2.5 (~1.24%) ou 2 (~5%). O limiar deve refletir o custo de falsos positivos vs falsos negativos no teu contexto. Por exemplo, numa operação com 1M de registos por dia, aplicar |z|>3 pode devolver ~2700 registos suspeitos por dia (pressuposição de normalidade), um volume gerível para revisão manual.
threshold = 3.0
anomalies = result.filter(abs(col("z_score")) > threshold)
anomalies.show()
Além da deteção, guarda a coluna z_score como característica para modelos (normaliza por grupo antes de treinar), ou calcula percentis por grupo para regras de negócio. Testa sempre com grupos de tamanhos variados e valida a taxa de alarmes com dados rotulados, se existirem.
Verificar o resultado
Confirma que as médias e desvios-padrão por grupo fazem sentido usando stats.show(). Verifica manualmente alguns cálculos (por exemplo, os exemplos numéricos dados) e testa grupos com um único registo para confirmar o tratamento de divisão por zero. Se os resultados forem estranhos, revê agregações e valores nulos.
Conclusão
Normalizar por grupo com Z-score em Apache Spark é uma técnica simples mas poderosa para deteção de anomalias e preparação de características. Para produção, incorpora este cálculo num pipeline ETL, valida limiares com dados históricos e escolhe stddev_pop vs stddev_samp conforme a definição estatística necessária. Próximo passo: aplicar por janelas temporais (ex.: Z-score por grupo e por mês) ou integrar numa etapa de validação automática.