Inteligência artificial, sem ruído.
Tutoriais10 min

PySpark Window Functions: guia prático para substituir o groupBy

Aprenda a usar window functions no PySpark para rankings, totais acumulados e médias móveis sem perder o detalhe das linhas originais.

PySpark Window Functions: guia prático para substituir o groupBy

Por que as window functions importam agora

Se você trabalha com dados em escala no PySpark, em algum momento vai esbarrar no limite do groupBy(). Ele é ótimo para agregar — somar mil linhas e devolver uma —, mas tem uma restrição fundamental: sempre retorna uma linha por grupo, apagando o detalhe das linhas originais.

As window functions resolvem exatamente isso. Elas calculam valores entre registros relacionados sem colapsar esses registros em um único resultado. Você mantém o detalhe transacional e, ao mesmo tempo, ganha informações de agregação do grupo — rankings, totais acumulados, comparações com registros anteriores e médias móveis. É a diferença entre “quanto vendeu cada loja” e “quanto vendeu cada loja e quanto cada venda representa do total”.

✅ O que você ganha

  • Ranking dentro de grupos — os maiores clientes de cada região, os melhores produtos por categoria.
  • Totais acumulados (running totals) que preservam a ordem cronológica dos eventos.
  • Comparação temporal com lag() e lead() — quanto mudou de um dia para o outro.
  • Participação relativa — o peso de cada transação no total do grupo, sem joins extras.
  • Médias móveis e janelas deslizantes para suavizar séries temporais.

⚠️ O que você NÃO ganha

  • Redução de linhas — window functions adicionam colunas, não resumem o DataFrame.
  • Desempenho de graça — janelas frequentemente exigem shuffle e ordenação, que custam caro em dados grandes.
  • Semântica temporal automática — uma janela de “7 linhas” não é “7 dias”; é preciso cuidado.

Requisitos

ComponenteMínimoRecomendadoIdeal
Python3.93.10+3.11+
PySpark3.33.53.5+ com Java 11/17
Memória4 GB8 GB16 GB+
Conhecimento prévioSQL básicogroupBy/aggDataFrames e Catalyst
Tempo estimado20–30 minutos para acompanhar os exemplos
Pré-requisitos para acompanhar este tutorial.

Passo 1 — Instalação com uv

Se o PySpark ainda não está instalado, o caminho mais rápido é o uv (funciona em Linux, macOS e Windows):

mkdir pyspark-windows
cd pyspark-windows
uv init
uv venv
source .venv/bin/activate   # no Windows: .venv\Scripts\activate
uv pip install pyspark

Em seguida, crie uma sessão Spark para testar a instalação. O parâmetro local[*] roda o Spark localmente usando todos os núcleos disponíveis — você não precisa de um cluster:

import os, sys
os.environ["PYSPARK_PYTHON"] = sys.executable
os.environ["PYSPARK_DRIVER_PYTHON"] = sys.executable
from pyspark.sql import SparkSession

spark = (SparkSession.builder
    .master("local[*]")
    .appName("sales-analysis")
    .config("spark.pyspark.python", sys.executable)
    .config("spark.pyspark.driver.python", sys.executable)
    .getOrCreate())

As duas primeiras linhas apontam o PySpark para o interpretador Python correto do seu ambiente virtual — sem isso, o Spark pode usar um Python errado e falhar silenciosamente.

Passo 2 — O que é uma janela (window)

Uma janela define o conjunto de linhas que o Spark considera ao calcular um valor para a linha atual. Ela tem três partes opcionais:

  • partitionBy() — divide os dados em grupos independentes.
  • orderBy() — define a ordem das linhas dentro de cada grupo.
  • rowsBetween() / rangeBetween() — define o “quadro” (frame) relativo à linha atual.
from pyspark.sql.window import Window
store_window = Window.partitionBy("store")

Isso divide os dados por loja. Cada transação de Londres pertence a uma partição, cada uma de Manchester a outra, e assim por diante. Aplicando uma agregação sobre essa janela:

from pyspark.sql import functions as F
sales_with_store_total = sales.withColumn(
    "store_total", F.sum("amount").over(store_window))

O resultado ainda contém cada transação, mas agora cada linha carrega o total de vendas da sua loja. Compare com o groupBy(), que devolveria apenas três linhas (uma por loja). A regra de ouro: use groupBy() quando quer uma linha de resultado por grupo; use window function quando quer cálculos de grupo ao lado das linhas originais.

Passo 3 — Ranking com row_number, rank e dense_rank

O ranking é um dos usos mais comuns. Para classificar transações da maior para a menor dentro de cada loja:

sales_rank_window = (Window
    .partitionBy("store")
    .orderBy(F.col("amount").desc()))

ranked_sales = sales.withColumn(
    "sale_rank", F.row_number().over(sales_rank_window))

O PySpark oferece três funções de ranking que diferem apenas quando há empates:

FunçãoComportamento com empatesExemplo (100, 100, 80)
row_number()Número sequencial único, sem empate1, 2, 3
rank()Empates recebem a mesma posição e deixam buracos1, 1, 3
dense_rank()Empates recebem a mesma posição sem buracos1, 1, 2
Diferenças entre as funções de ranking.

Use row_number() quando precisa de exatamente um primeiro, segundo e terceiro lugar. Use rank() ou dense_rank() quando empates devem receber tratamento igual.

Um dos padrões mais úteis é selecionar os N maiores registros de cada grupo:

top_two_sales_per_store = (sales
    .withColumn("sale_rank", F.row_number().over(sales_rank_window))
    .filter(F.col("sale_rank") <= 2)
    .orderBy("store", "sale_rank"))

Atenção: isso é diferente de sales.orderBy(F.col("amount").desc()).limit(2), que devolve os dois maiores de todo o conjunto. A versão com janela devolve os dois maiores de cada loja.

Passo 4 — Totais acumulados (running totals)

Um total acumulado soma o valor da linha atual aos valores das linhas anteriores. Para calcular vendas cumulativas por loja, defina uma janela que particiona por loja, ordena por data e começa na primeira linha da partição até a linha atual:

running_total_window = (Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id")
    .rowsBetween(Window.unboundedPreceding, Window.currentRow))

sales_with_running_total = sales.withColumn(
    "running_store_total", F.sum("amount").over(running_total_window))

A segunda coluna de ordenação (transaction_id) torna a ordem determinística quando duas transações compartilham a mesma data — sem um desempate claro, linhas com valores iguais podem aparecer em ordem imprevisível.

Passo 5 — Comparando com a linha anterior (lag e lead)

A função lag() recupera o valor de uma linha anterior na mesma janela. É essencial para medir mudanças ao longo do tempo:

store_date_window = (Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id"))

sales_with_previous_amount = (sales
    .withColumn("previous_amount", F.lag("amount").over(store_date_window))
    .withColumn("change_from_previous",
        F.col("amount") - F.col("previous_amount")))

A primeira transação de cada loja não tem antecessora, então previous_amount fica NULL. A função lead() faz o mesmo, mas olhando para a frente (a próxima linha). Ambas ajudam a identificar lacunas na atividade de clientes, eventos atrasados ou mudanças em medições diárias.

Passo 6 — Participação de cada linha no total do grupo

Agregações de janela facilitam comparar cada linha com o seu grupo. Para calcular o percentual de cada transação no total da loja:

store_total_window = Window.partitionBy("store")
sales_with_share = (sales
    .withColumn("store_total", F.sum("amount").over(store_total_window))
    .withColumn("share_of_store_sales",
        F.col("amount") / F.col("store_total")))

Como o total da loja aparece ao lado de cada transação, não é preciso agregar e depois fazer join de volta às linhas originais. O mesmo padrão calcula o salário de um funcionário como proporção da folha do departamento, ou as vendas de um produto como proporção da sua categoria.

Passo 7 — Médias móveis e janelas deslizantes

Um total acumulado inclui todas as linhas anteriores. Uma média móvel usa um número limitado de linhas vizinhas. Para a média das três últimas transações:

moving_average_window = (Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id")
    .rowsBetween(-2, Window.currentRow))

sales_with_moving_average = sales.withColumn(
    "three_sale_average", F.avg("amount").over(moving_average_window))

Na primeira linha de cada loja, a média usa uma transação; na segunda, duas; da terceira em diante, a transação atual mais as duas anteriores. Note que isso é uma janela baseada em linhas, não em tempo: se uma loja vende várias vezes por dia e outra vende uma vez por semana, cada cálculo ainda cobre três linhas.

Passo 8 — rowsBetween vs rangeBetween

Os frames controlam quais linhas contribuem para o cálculo:

  • rowsBetween() usa posições de linha. .rowsBetween(-3, Window.currentRow) inclui a linha atual e as três anteriores.
  • rangeBetween() usa valores da coluna de ordenação. Se a coluna contém timestamps em segundos, uma janela de 7 dias ficaria assim:
seconds_in_seven_days = 7 * 24 * 60 * 60
seven_day_window = (Window
    .partitionBy("store")
    .orderBy(F.col("sale_timestamp").cast("long"))
    .rangeBetween(-seconds_in_seven_days, 0))

A distinção importa: use rowsBetween() quando quer um número fixo de registros; use rangeBetween() quando quer registros dentro de um intervalo de valor ou tempo. Janelas baseadas em tempo exigem cuidado — a expressão de ordenação precisa usar uma representação numérica adequada com unidades consistentes.

Passo 9 — Reutilizando especificações de janela

Especificações de janela não modificam o DataFrame sozinhas — elas descrevem como agrupar, ordenar e enquadrar. Definir uma vez e reutilizar torna o código mais legível:

store_total_window = Window.partitionBy("store")
store_date_window = (Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id"))
running_total_window = (store_date_window
    .rowsBetween(Window.unboundedPreceding, Window.currentRow))

Considerações de desempenho

Window functions são úteis, mas não são de graça. O PySpark pode precisar mover e ordenar dados para que linhas com a mesma chave de partição sejam processadas juntas — esse shuffle tem custo. Para inspecionar o plano de execução:

analysed_sales.explain("formatted")

Procure por operações de exchange e sort. Quatro hábitos ajudam a manter as consultas sob controle:

  • Filtre cedo — remova linhas desnecessárias antes de aplicar a janela.
  • Selecione só as colunas necessárias antes da operação de janela.
  • Observe partições distorcidas — se uma chave concentra muito mais linhas que as outras, uma tarefa terá muito mais trabalho.
  • Reuse resultados com cuidado.cache() só compensa quando o resultado é reutilizado várias vezes.

❌ Erros comuns (e como corrigir)

  1. Esquecer o partitionBy()Window.orderBy(...) ranqueia todas as linhas do conjunto inteiro, não por grupo. Adicione .partitionBy("coluna") quando a intenção é agrupar.
  2. Ordenação incompleta — se várias linhas têm a mesma data, ordenar só por data deixa a ordem ambígua. Adicione um desempate como transaction_id.
  3. Esperar que a janela reduza linhas — window functions adicionam colunas, não resumem. Para manter só os mais bem ranqueados, adicione a coluna de rank e filtre depois.
  4. Tratar janela de linhas como janela de temporowsBetween(-6, 0) inclui sete linhas, não sete dias. Use rangeBetween() para períodos reais.
  5. Coluna NULL no lag() quebrando o cálculo — a primeira linha de cada partição retorna NULL; trate com F.coalesce() ou filtre os NULL antes de usar em aritmética.

❓ Perguntas frequentes

  1. Quando usar groupBy() em vez de window? Quando você precisa de uma única linha de resultado por grupo e não se importa em perder o detalhe das linhas originais.
  2. Preciso de um cluster Spark para testar? Não. O modo local[*] roda tudo na sua máquina, usando os núcleos disponíveis.
  3. Qual a diferença entre rank() e dense_rank()? Ambos dão a mesma posição para empates, mas rank() pula números depois do empate (1, 1, 3) enquanto dense_rank() não pula (1, 1, 2).
  4. Window functions funcionam com DataFrames e SQL? Sim. Você pode usar as funções equivalentes em SQL via spark.sql("SELECT ..., row_number() OVER (PARTITION BY store ORDER BY amount DESC) ...").
  5. Por que minha consulta ficou lenta depois de adicionar uma janela? Provavelmente o shuffle e a ordenação. Filtre cedo, selecione menos colunas e inspecione o plano com .explain("formatted").

Para onde isso vai

À medida que os pipelines de dados crescem e as equipes migram de groupBy() para análises mais granulares, dominar window functions deixa de ser diferencial e vira pré-requisito. Com o Spark consolidado como padrão de processamento distribuído e a crescente adoção de lakehouse e streaming estruturado, entender como partição, ordenação e frames se combinam é uma das habilidades mais transferíveis da engenharia de dados — e vai continuar valendo quando os volumes de dados multiplicarem.


Descubra mais sobre noticiAI

Assine para receber nossas notícias mais recentes por e-mail.

R
Sobre o autorRedação Noticiai

Equipe editorial dedicada a explicar inteligência artificial com clareza, independência e contexto.