Pular para o conteúdo principal

Dominando o groupBy do PySpark para agregação de dados em escala

Explore o método groupBy do PySpark, que permite a profissionais de dados aplicar funções de agregação nos seus dados. É uma forma poderosa de particionar e resumir rapidamente grandes conjuntos de dados, aproveitando as técnicas avançadas do Spark.
Atualizado 17 de set. de 2026  · 15 min lido

Explorar com IA

ChatGPTClaudePerplexity

No universo de big data, raramente olhamos os dados em bruto. Para tornar as informações mais digeríveis, agrupar e agregar dados é uma operação comum e poderosa. Ao trabalhar com sistemas distribuídos como Apache Spark e Python, contar com a função groupBy do PySpark permite reunir e resumir dados distribuídos. Ela espelha a funcionalidade da cláusula GROUP BY do SQL, mas foi criada para lidar com processamento distribuído em conjuntos de dados massivos com eficiência.

O groupBy do PySpark permite particionar os dados com base em várias colunas, que depois podem ser agregadas em diferentes medidas, como somas, médias e assim por diante. Para fazer isso de forma eficiente, o PySpark segue o paradigma split-apply-combine:

  1. Divide os dados em grupos com base em um critério,
  2. Aplica a lógica de agregação ou transformação a cada grupo,
  3. Combina os resultados em um novo DataFrame.

Gráfico mostrando como o método groupBy do PySpark funciona: divide por chave, aplica soma a cada chave e combina os dados agregados

Exemplo do método groupBy() do PySpark aplicando soma

O PySpark usa avaliação preguiçosa (lazy evaluation), o que significa que operações como groupBy só são computadas quando ações (como show() ou collect()) são chamadas. Antes disso, o PySpark constrói um DAG (Directed Acyclic Graph) inicial que é refinado para performance. Isso ajuda o Spark a otimizar o plano de consulta antes da execução.

Se você ainda não teve oportunidade de trabalhar com PySpark, vale conferir os Fundamentos de Big Data com PySpark.

Criando um DataFrame no PySpark

Primeiro de tudo, precisamos de um DataFrame do PySpark. Se quiser relembrar os principais comandos, dê uma olhada neste ótimo cheatsheet de DataFrame do PySpark.

Vamos passar por etapas e comandos essenciais para começar com o PySpark. Antes de iniciar, garanta que o PySpark esteja instalado no seu ambiente Python, assim como Java e o Java JDK. Se precisar de ajuda, siga este guia de como começar com PySpark.

Inicializando a SparkSession

O primeiro passo é colocar o PySpark para rodar!

from pyspark.sql import SparkSession

# Start a SparkSession
spark = SparkSession.builder \
    .appName("GroupByExample") \
    .getOrCreate()

DataFrame de exemplo para GroupBy

Vamos construir um DataFrame de exemplo para trabalharmos.

from pyspark.sql import Row

data = [
    Row(department="Sales", employee="Alice", salary=5000),
    Row(department="Sales", employee="Bob", salary=4800),
    Row(department="HR", employee="Carol", salary=4000),
    Row(department="HR", employee="David", salary=3900),
    Row(department="IT", employee="Eve", salary=6000)
]

df = spark.createDataFrame(data)
df.show()

Saída:

	+----------+--------+------+
	|department|employee|salary|
	+----------+--------+------+
	|     Sales|   Alice|  5000|
	|     Sales|     Bob|  4800|
	|        HR|   Carol|  4000|
	|        HR|   David|  3900|
	|        IT|     Eve|  6000|
	+----------+--------+------+

O que é o GroupBy do PySpark?

O groupBy do PySpark é uma transformação usada para dividir os dados em grupos com base em uma ou mais colunas, que depois podem ser agregados ou transformados de forma independente.

Sintaxe básica e exemplos

Usar o método groupBy é simples: você o chama no dataframe de interesse. Dá para usá-lo em qualquer tipo de coluna, desde que faça sentido particionar por ela. Por exemplo, convém evitar floats.

grouped = df.groupBy("department")

Isso cria um objeto GroupedData, que não é um DataFrame em si. Ele orienta o Spark sobre como começar a particionar os dados. A partir desse GroupedData, você aplica agregações como .count() ou .sum() para obter um novo DataFrame.

grouped.count().show()

Parâmetros e valores de retorno

O único parâmetro do método é *cols, que aceita nomes de coluna, expressões de coluna, ordinais de coluna (int) ou uma lista de colunas. Se a coluna for algo pelo qual você deseja particionar, pode usá-la na operação. O retorno é sempre um objeto GroupedData.

# agrupando por nome de coluna, como acima
	df.groupBy("department")

# agrupando por expressão de coluna
	df.groupBy(df.department)

# agrupando por ordinal da coluna
	df.groupBy(1)

# agrupando por lista de colunas; você pode misturar os métodos!
	df.groupBy(["department", 2])

GroupBy em uma ou várias colunas

Como mostrado acima, você pode agrupar por uma ou várias colunas. Uma única coluna é ótima quando você se interessa por um único eixo, como departamentos ou anos. Já várias colunas servem para trazer mais camadas de detalhe, por exemplo, vendedores específicos em diferentes departamentos ou meses específicos dentro de um ano.

# agrupando por coluna única
	df.groupBy("department").sum(“salary”).show()

# agrupando por múltiplas colunas
	df.groupBy(["department", 2]).sum(“salary”).show()

Como você provavelmente notou, o método groupBy funciona de forma muito semelhante ao GROUP BY do SQL. Mais adiante, vamos mostrar como usar essa linguagem também no PySpark.

Funções e técnicas de agregação

O PySpark oferece uma ampla variedade de métodos de agregação nativos para trabalhar com dados agrupados. Se você já conhece SQL, vai reconhecer muitos deles, como count(), sum(), avg() e por aí vai. Você pode revisar nosso guia de funções de agregação em SQL.

Funções de agregação nativas

Vamos passar pelas funções de agregação nativas:

  • count(): conta o número de registros na partição
  • sum(): soma os valores numéricos
  • avg(): calcula a média dos valores numéricos
  • min(): retorna o menor valor da partição
  • max(): retorna o maior valor da partição

Além das funções de agregação, você pode anexar o método .alias() para renomear colunas e deixá-las mais fáceis de entender. Vamos mostrar um exemplo abaixo.

Múltiplas agregações com agg()

Em vez de executar cada agregação separadamente, você pode fazer várias de uma vez usando o método agg(). Cada agregação cria uma nova coluna no seu dataframe. Isso também reduz a necessidade de várias chamadas a groupBy, o que melhora a performance. Não é preciso agregar na mesma coluna com agg(); você pode definir uma coluna diferente para cada função de agregação.

# Depois de iniciar a sessão
from pyspark.sql import functions as sf

df.groupBy("department").agg(
    sf.count("employee").alias("employee_count"), 
    sf.avg("salary").alias("avg_salary"),
    sf.max("salary").alias("max_salary")
).show()

Saída:

+----------+--------------+----------+----------+
|department|employee_count|avg_salary|max_salary|
+----------+--------------+----------+----------+
|     Sales|             2|    4900.0|      5000|
|        HR|             2|    3950.0|      4000|
|        IT|             1|    6000.0|      6000|
+----------+--------------+----------+----------+

Padrões avançados de agregação

Uma das grandes vantagens do PySpark é poder usar padrões de agregação mais avançados, como pivotar os dados, fazer rollup e criar cubos de dados. Você também pode criar grouping sets.

Pivot

Assim como criar uma tabela dinâmica no Excel, você pode pivotar os dados em diferentes colunas. Neste caso, estamos agrupando por department e pivotando a coluna employee para ver o salário total de cada pessoa.

Isso significa que cada linha mostra o departamento e cada coluna representa os funcionários daquele departamento. Dá para ver como isso seria poderoso se tivéssemos dados por ano e quiséssemos comparar o desempenho de cada departamento ano a ano.

df.groupBy("department").pivot("employee").sum("salary").show()

Rollups e cubes

Duas formas bem poderosas de agregar dados são rollup() e cube(). Enquanto um groupby() simples mostra os resultados para as agregações existentes, rollup() e cube() são estruturas hierárquicas. Ou seja, agregam em níveis mais granulares.

Por exemplo, rollup() agrega da esquerda para a direita, mostrando cada permutação possível se fôssemos iterar passo a passo. Para cada departamento, mostramos cada pessoa do departamento e também o grupo final (agregado) daqueles que não estão especificados.

Já o cube() mostra todas as permutações possíveis para todas as colunas agregadas: cada departamento separadamente, cada funcionário separadamente e todas as combinações entre ambos.

Resumindo:

  • Rollup cria subtotais hierárquicos seguindo a ordem das colunas.
  • Cube gera subtotais para todas as combinações possíveis das colunas especificadas.

Veja um exemplo no código e na saída abaixo:

rollup() Código:

# usando rollup
df.rollup("department", “employee”).sum("salary").show()

rollup() Saída:

+----------+--------+-----------+
|department|employee|sum(salary)|
+----------+--------+-----------+
|      NULL|    NULL|      23700|
|     Sales|   Alice|       5000|
|     Sales|    NULL|       9800|
|     Sales|     Bob|       4800|
|        HR|   Carol|       4000|
|        HR|    NULL|       7900|
|        HR|   David|       3900|
|        IT|     Eve|       6000|
|        IT|    NULL|       6000|
+----------+--------+-----------+

cube() Código:

# usando cube
df.cube("department", "employee").sum("salary").show()

cube() Saída:

+----------+--------+-----------+
|department|employee|sum(salary)|
+----------+--------+-----------+
|      NULL|   Alice|       5000|
|      NULL|    NULL|      23700|
|     Sales|   Alice|       5000|
|     Sales|    NULL|       9800|
|     Sales|     Bob|       4800|
|      NULL|     Bob|       4800|
|      NULL|   Carol|       4000|
|        HR|   Carol|       4000|
|        HR|    NULL|       7900|
|        HR|   David|       3900|
|      NULL|   David|       3900|
|        IT|     Eve|       6000|
|        IT|    NULL|       6000|
|      NULL|     Eve|       6000|
+----------+--------+-----------+

Grouping sets

Grouping sets permitem definir múltiplos níveis de agregação. No exemplo, estou agregando no nível de departamento e funcionário, apenas por departamento e o total geral para obter o salário total.

Você vai notar que a sintaxe do método groupingSets() é um pouco diferente. Primeiro você define a lista de conjuntos [(“department”, “employee”), (“department”, ), ()], em que o primeiro é departamento e funcionário, o segundo apenas departamentos e o último () significa todos (total geral).

Depois você define as colunas de agregação dentro dos conjuntos. O restante é escrito normalmente usando .agg().

df.groupingSets([("department", "employee"), ("department",), ()], "department","employee").agg(sf.sum("salary")).sort("department","employee").show()

Saída:

+----------+--------+-----------+
|department|employee|sum(salary)|
+----------+--------+-----------+
|      NULL|    NULL|      23700|
|        HR|    NULL|       7900|
|        HR|   Carol|       4000|
|        HR|   David|       3900|
|        IT|    NULL|       6000|
|        IT|     Eve|       6000|
|     Sales|    NULL|       9800|
|     Sales|   Alice|       5000|
|     Sales|     Bob|       4800|
+----------+--------+-----------+

Funções de agregação personalizadas

Por fim, podemos criar funções de agregação personalizadas usando User Defined Functions (UDFs) no PySpark. Existem duas formas: udf e pandas_udf. Ambas permitem criar funções customizadas, cada uma com seus prós e contras.

A udf padrão do Spark permite criar funções nativas do Spark. Isso significa que a sintaxe e os tipos de dados precisam funcionar nativamente no Spark. Embora isso limite um pouco o que pode ser feito na UDF, você aproveita ao máximo o poder computacional distribuído do Spark, sendo melhor para conjuntos de dados maiores.

from pyspark.sql.functions import udf
from pyspark.sql.types import IntegerType

@udf(returnType=IntegerType())
def bonus(salary):
    return int(salary * 0.1)

df.withColumn("bonus", bonus(df.salary)).show()

Já a pandas_udf permite criar funções personalizadas mais "pythônicas". Você fica limitado apenas ao que é possível no Pandas, e não só no Spark.

Porém, isso significa que você não aproveita o poder computacional distribuído do Spark e passa a depender de computação local. É mais indicada para conjuntos de dados pequenos a médios.

import pandas as pd
from pyspark.sql.functions import pandas_udf

@pandas_udf("double")
def salary_bonus(s: pd.Series) -> pd.Series:
    return s * 0.1

df.withColumn("bonus", salary_bonus(df.salary)).show()

No geral, use funções personalizadas apenas quando não houver uma forma melhor de realizar a agregação. Não reinvente a roda!

Além disso, se algo puder ser vetorizado (como as funções acima), prefira uma versão vetorizada em vez de uma UDF de agregação. Para saber mais sobre groupBy(), leia este artigo que aprofunda o framework split-apply-combine com pandas.

Filtrando dados agregados

No PySpark, você pode filtrar grupos com base em métricas agregadas após o agrupamento usando o método filter(). Nele, você fornece uma condição em expressões Python ou SQL.

Aliás, você também pode usar where() se preferir, pois é apenas um alias de filter() e realiza a mesma operação.

Você pode escolher filtrar antes ou depois da agregação. Filtrar antes impacta a agregação ao limitar quais dados serão agregados e pode melhorar a performance.

Por exemplo, talvez a gente só queira contar o número de funcionários acima de um certo salário para encontrar o número de colaboradores com "salário alto".

# filtrar por salários altos
filter_df = df.filter(df.salary > 4000)

# agregar e encontrar o total de funcionários com salário alto
agg_filter_df = filter_df.groupBy("department").agg(sf.count("*").alias("high_salary_emp"))
agg_filter_df.show()

Já filtrar depois da agregação não impacta a agregação original. Talvez a gente queira contar todos os funcionários de um departamento e só visualizar departamentos acima de um certo tamanho. Isso adiciona uma camada extra de processamento, mas não deve aumentar muito a carga computacional.

# criando um dataframe com a contagem de funcionários
agg_df = df.groupBy("department").agg(sf.count("*").alias("num_employees"))

# filtrando departamentos com mais de 1 funcionário
agg_df.where("num_employees > 1").show()

Escolha bem quando filtrar, pois isso impacta a precisão e o resultado final. Filtrar antes pode reduzir o total final, enquanto filtrar depois pode deixar muitos no resultado.

Estratégias de otimização de performance do groupBy no PySpark

Embora o PySpark otimize automaticamente a operação groupBy, existem estratégias para aumentar ainda mais a velocidade de processamento. Coisas como minimizar shuffles de dados, mitigar skew e otimizar a execução contribuem para melhorar a performance do groupBy.

Gerenciamento de shuffle

O groupBy causa um shuffle, que redistribui dados entre partições para agrupar chaves semelhantes. Shuffles consomem rede e disco, então minimizá-los ou otimizá-los é essencial. Algumas técnicas importantes:

  • Use repartition() com critério: você pode indicar por qual coluna os dados devem ser particionados, reduzindo o escopo de busca.
  • Ajuste spark.sql.shuffle.partitions. Por padrão, o PySpark usa 200 partições de shuffle. Se o dataset for menor, você pode reduzir esse número.
  • Habilite compressão no shuffle. Comprimir reduz o overhead de rede e disco.

Veja alguns exemplos de como melhorar o gerenciamento de shuffle.

# reduzir o número de partições de shuffle
spark.conf.set("spark.sql.shuffle.partitions", "64")  # Ajuste conforme o tamanho do cluster

# garantir que o spark comprima os dados
spark.conf.set("spark.shuffle.compress", "true") # comprime a transferência de rede
spark.conf.set("spark.shuffle.spill.compress", "true") # comprime o spill em disco

# reparticiona os dados antes do agrupamento para otimizar a agregação
df.repartition("department").groupBy("department").sum("salary").show()

Técnicas para mitigar skew

O skew acontece quando certas chaves aparecem muito mais do que outras, gerando cargas desiguais nas partições e tarefas lentas. Isso sobrecarrega alguns workers mais do que outros. Podemos minimizar o skew e distribuir o trabalho de forma mais uniforme com salting, evitando joins enviesados e usando broadcast joins para dimensões menores.

  • Salting: adicione uma coluna de números aleatórios para forçar distribuição uniforme entre os workers
  • Reparticionar para minimizar skew: forçar o Spark a usar uma coluna diferente para particionar pode melhorar a distribuição de carga
# Exemplo de salting:

df = df.withColumn("salted_key", sf.rand()) # cria coluna aleatória
df = df.repartition(2, 'salted_key') # usa para reparticionar os dados
df.groupBy(sf.spark_partition_id()).count().show()

Otimização de execução

Os planos lógico e físico de execução dos jobs PySpark podem ser otimizados com recursos nativos. Eles melhoram a performance das suas agregações com groupBy.

  • Otimização Catalyst: o Spark reescreve planos ineficientes automaticamente; escrever transformações declarativas (não laços procedurais) ajuda o otimizador.
  • Cache: fazer cache é útil quando o mesmo resultado de groupBy é usado várias vezes no pipeline.
  • Broadcast joins: faça broadcast de dataframes menores para mantê-los em memória enquanto os maiores são particionados, minimizando overhead de rede e disco no cluster.
grouped_df = df.groupBy("department").sum("salary").cache()
grouped_df.show() # você pode fazer cache de resultados intermediários

# use broadcasting (pseudocódigo aqui):
from pyspark.sql.functions import broadcast
df.join(broadcast(smaller_df), "department").show()

Para mais detalhes sobre como otimizar operações no PySpark, leia este artigo sobre joins no PySpark para entender a mecânica por trás.

Análise comparativa com operações em RDD

Embora a API de DataFrame seja a interface recomendada para a maioria dos usuários de PySpark, graças à sua abstração de alto nível e eficiência (via Catalyst Optimizer), entender a camada de RDD (Resilient Distributed Dataset) pode ser valioso, especialmente para quem quer mais controle ou está migrando código legado do Spark. Em geral, porém, usar a API de DataFrame ou SQL tende a entregar melhor performance.

O RDD é o componente central do PySpark, sobre o qual tudo (incluindo a API de DataFrame) é construído. Ele costuma trabalhar em memória e lida melhor com dados de streaming do que a API de DataFrame. Entretanto, os dados são imutáveis, então não podem ser alterados após a criação do RDD.

Esta seção compara o groupBy na API de DataFrame com seus equivalentes na API de RDD, analisando performance e aderência ao caso de uso.

Primeiro, vamos transformar o dataframe em um RDD.

# Converter DataFrame para RDD de (chave, valor)
rdd = df[['department','salary']].rdd
rdd.collect()

Agora podemos executar operações como groupByKey() e reduceByKey() no RDD. Vamos ver o groupByKey() primeiro.

O método groupByKey() envia todos os valores com a mesma chave para o mesmo executor e depois os agrupa. Por isso, pode levar bastante tempo para redistribuir os dados para a chave correta. Isso também pode gerar skew.

rdd.groupByKey().mapValues(list).collect()

Em vez disso, podemos usar o método reduceByKey(). Ele primeiro agrega em cada partição e depois envia os dados pela rede. Isso minimiza a quantidade de dados em shuffle e, como resultado, funciona mais rápido.

Portanto, reduceByKey é fortemente recomendado em vez de groupByKey porque reduz significativamente as operações de shuffle.

Na prática, ele agrega em cada partição, depois combina os dados e então agrega novamente. É ideal para agregações que unem chaves, como sum() e max().

from operator import add

rdd.reduceByKey(add).collect()

Benefícios de performance: DataFrame vs métodos de agregação em RDD

Aqui vai uma tabela que resume as vantagens de cada método de agregação.

Recurso

DataFrame groupBy

RDD groupByKey()

RDD reduceByKey()

Nível de abstração

Alto

Baixo

Baixo

Execução otimizada

Sim (Catalyst & Tungsten)

Não

Não

Minimização de shuffle

Sim

❌ (shuffle completo)

✅ (com combiner)

Eficiência de memória

Alta

Baixa

Média–Alta

Flexibilidade

Moderada

Alta

Média

Performance em Big Data

Excelente

Fraca

Boa

Uso recomendado

Maioria dos casos

Raro/especializado

Agregações customizadas e de baixo nível

O curso Fundamentos de Big Data com PySpark cobre programação com RDDs em mais detalhes e descreve como eles são a espinha dorsal do PySpark.

Consulta SQL GROUP BY no PySpark

Outra maneira de fazer agregações no PySpark é usar a API SQL para escrever instruções em SQL. É uma ótima opção se você se sente mais confortável com SQL.

O primeiro passo é criar uma visão temporária com o método createOrReplaceTempView() do DataFrame. Depois, use spark.sql() para escrever sua consulta.

# Criar uma visão temporária a partir do DataFrame
df.createOrReplaceTempView("employees")

# Escrever uma instrução em SQL
spark.sql("""
    SELECT department, AVG(salary) AS avg_salary
    FROM employees
    GROUP BY department
""").show()

Saída:

+----------+----------+
|department|avg_salary|
+----------+----------+
|     Sales|    4900.0|
|        HR|    3950.0|
|        IT|    6000.0|
+----------+----------+

Como você vê, é tão fácil quanto escrever uma query SQL normal. Os principais motivos para usar a API SQL são familiaridade e facilidade para escrever consultas. A API de DataFrame costuma ser melhor por fornecer tipos de dados consistentes e, geralmente, otimizar mais.

Você também ganha acesso a métodos de DataFrame que não teria com a API Spark SQL. Para saber mais sobre o GROUP BY no SQL, leia este artigo sobre GROUP BY e HAVING em SQL.

Aplicações no mundo real

Existem muitas aplicações práticas para agregar seus dados. Na verdade, é quase certo que você vai precisar agregar de alguma forma para analisar e compartilhar resultados.

Abaixo, trazemos alguns exemplos de diferentes casos de uso de agregações. Alguns podem ser pseudocódigo ou não se aplicar exatamente ao nosso dataset de teste; a ideia é inspirar como você pode escrever essas agregações.

Business intelligence

Um tema recorrente que discutimos é o uso de groupBy como meio de análise hierárquica. Imagine que temos uma coluna "revenue" e queremos ver o desempenho de cada departamento. Podemos agrupar por departamento e avaliar a contribuição de receita de cada um.

CopyEdit
df.groupBy("department").agg(sum("revenue").alias("total_revenue")).show()

Análise de séries temporais

Agregações dão acesso a funções de janela (window functions). Essas funções analisam uma janela deslizante de dados, seja por linhas sequenciais ou por períodos de tempo sequenciais.

Primeiro definimos uma janela usando o objeto Window para indicar o que queremos partitonBy e orderBy. Em seguida, usamos esse objeto no método over() da função de agregação escolhida. Por exemplo, uma média móvel ficaria sf.avg().over(window).show().

Imagine que queremos ver a média móvel de salários por departamento. Agruparíamos por department e ordenaríamos por uma nova coluna que podemos chamar de employeeId. Como IDs de funcionários costumam ser sequenciais, dá para ver como a média salarial mudou ao longo do tempo.

from pyspark.sql.window import Window

# Definir a janela que queremos particionar e ordenar
windowSpec = Window.partitionBy("department").orderBy("employeeId")

# Executar sf.avg na coluna “salary” sobre a windowSpec
df.withColumn("rolling_avg", sf.avg("salary").over(windowSpec)).show()

Isso é bem parecido com como fazemos funções de janela em SQL.

Analytics para mídia

Talvez queiramos entender várias métricas de usuário em uma empresa de mídia. Queremos, por exemplo, tempo total assistido e número de vídeos únicos. Podemos fazer groupBy() em uma coluna user_id e depois .agg() com funções como sum() e countDistinct() para diferentes métricas.

df.groupBy("user_id").agg(
    sf.sum("watch_time").alias("total_watch"), # soma de minutos assistidos
    sf.countDistinct("video_id").alias("unique_views") # contagem de vídeos únicos assistidos
).show()

Boas práticas de GroupBy no PySpark e como evitar armadilhas

Ao agregar dados no PySpark, é fácil cair em armadilhas que levam a baixa performance e tempos longos de execução. Seguir estas boas práticas ajuda a garantir performance, correção e escalabilidade.

Armadilhas comuns

Aqui estão alguns jeitos comuns de transformar consultas simples em programas que demoram a rodar.

  1. Uso excessivo de groupByKey() em RDDs:
    • Evite groupByKey() em RDDs, pois seu mecanismo faz o PySpark distribuir dados entre partições. Isso pode gerar grande overhead de rede e disco.
    • Em vez disso, prefira reduceByKey() ou fique na API de DataFrame.
  2. Não tratar dados com skew:
    • Quando um grupo domina (por exemplo, um departamento com milhões de registros), essa tarefa vira gargalo.
    • Use salting ou particionamento customizado para mitigar, redistribuindo a carga com repartition().
  3. Falhar em encadear funções de agregação:
    • Várias linhas separadas com funções de agregação — por exemplo, sum() em um groupBy e depois count() no mesmo conjunto agrupado em outra linha — quebram a capacidade de otimização do PySpark e geram planos ineficientes.
    • O PySpark usa avaliação preguiçosa, ou seja, gera todo o plano antes das ações. Encadear transformações (ex.: groupBy().agg().filter()) permite ao Catalyst otimizar melhor.
  4. Uso incorreto de UDFs:
    • UDFs em Python impedem o Spark de otimizar completamente as consultas, pois o Catalyst não consegue otimizá-las.
    • Se possível, use funções nativas do Spark ou pandas UDFs para melhor performance.
  5. Mau gerenciamento de memória:
    • Agregações podem ser custosas e consumir muita memória. Isso pode levar a exaustão quando você reutiliza resultados de agregação várias vezes.
    • Monitore o uso de memória e considere usar persist() ou cache() ao reutilizar dados agrupados.

Dicas de otimização e eficiência

Aqui vão algumas formas de otimizar suas agregações no PySpark para manter tudo fluindo bem!

Tratamento de null

  • Agregar sobre colunas com null pode gerar resultados inesperados.
  • Use na.fill() ou na.drop() antes da agregação.
df.na.fill({"salary": 0}).groupBy("department").sum("salary").show()

Evite shuffles excessivos

  • Como discutido, shuffles deixam tudo mais lento. Reparticione logicamente antes de agrupar para minimizar o tamanho do shuffle.
  • Ajuste spark.sql.shuffle.partitions com base no volume de dados e teste o melhor número de partições.

Use planos de explicação

  • Monitore os planos com .explain() para entender a execução física.
  • Procure por sinais de shuffles amplos, hints de broadcast e varreduras ineficientes.
df.groupBy("department").sum("salary").explain(True)

Saída:

	== Parsed Logical Plan ==
	'Aggregate ['department], ['department, unresolvedalias('sum(salary#600L))]
	+- LogicalRDD [department#598, employee#599, salary#600L], false

	== Analyzed Logical Plan ==
	department: string, sum(salary): bigint
	Aggregate [department#598], [department#598, sum(salary#600L) AS sum(salary)#679L]
	+- LogicalRDD [department#598, employee#599, salary#600L], false

	== Optimized Logical Plan ==
	Aggregate [department#598], [department#598, sum(salary#600L) AS sum(salary)#679L]
	+- Project [department#598, salary#600L]
	   +- LogicalRDD [department#598, employee#599, salary#600L], false

	== Physical Plan ==
	AdaptiveSparkPlan isFinalPlan=false
	+- HashAggregate(keys=[department#598], functions=[sum(salary#600L)], output=[department#598, sum(salary)#679L])
	   +- Exchange hashpartitioning(department#598, 200), ENSURE_REQUIREMENTS, [plan_id=1841]
	      +- HashAggregate(keys=[department#598], functions=[partial_sum(salary#600L)], output=[department#598, sum#681L])
	         +- Project [department#598, salary#600L]
	            +- Scan ExistingRDD[department#598,employee#599,salary#600L]
	

Monitore performance com o Spark UI

  • Use o Spark Web UI para acompanhar estágios, tarefas e identificar agregações lentas ou skew.
  • Você pode saber mais sobre o Spark Web UI neste curso de Introdução ao Spark SQL em Python.

Aproveite o tuning de configuração

O Spark tem diversas configurações que podem ser ajustadas para workloads grandes de groupBy. Veja a tabela abaixo com algumas chaves importantes.

Configuração

Descrição

spark.sql.shuffle.partitions

Controla o número de partições para shuffles. Reduza para jobs pequenos, aumente para grandes. (Padrão: 200)

spark.sql.autoBroadcastJoinThreshold

Habilita broadcast automático de tabelas pequenas. Defina -1 para desabilitar ou aumente para suportar joins maiores. (Padrão: 10MB)

spark.executor.memory

Controla a memória disponível por executor. Aumente para grandes agregações. (Padrão: 4g)

spark.sql.adaptive.enabled

Habilita o Adaptive Query Execution (AQE), que pode otimizar dinamicamente shuffles, skew e joins. (Padrão: true)

Notas e detalhes avançados de implementação

Para usuários avançados ou quem lida com datasets muito grandes, entender como o groupBy se comporta por baixo dos panos é essencial. Vamos cobrir alguns pontos finos de como esse método funciona.

Avaliação preguiçosa e plano de execução

Todas as operações de DataFrame no PySpark, incluindo groupBy, são avaliadas de forma preguiçosa. Ou seja, nenhuma computação real acontece até que uma ação (como show() ou collect()) seja acionada.

Isso permite ao Catalyst reorganizar, combinar ou remover operações para melhorar a performance. Aproveite isso e encadeie múltiplas agregações e métodos para que o Spark crie um plano otimizado.

Tipo de retorno e colisões de nomes

Ao usar groupBy(), é retornado um objeto GroupedData. Executar uma função de agregação, como sum(), retorna então um objeto tipo DataFrame. Usando show(), você verá os resultados.

Lembre que groupBy().agg() retorna um novo DataFrame com colunas recém-nomeadas. Sempre use aliases com o método alias() para evitar colisões de nomes e deixar as saídas mais claras. Se você não usar aliases, isso pode levar a colisões de nomes de coluna ou problemas em joins mais adiante.

from pyspark.sql import functions as F

df.groupBy("department") \

  .agg(

      sum("salary").alias("total_salary"),

      count("*").alias("employee_count")

  )

Formato do resultado da agregação

Um ponto importante no PySpark: a agregação retornada não preserva a ordem das linhas. Use orderBy() para saídas previsíveis e ordenação consistente.

df.groupBy("department").sum("salary").orderBy("department").show()

Comportamentos por versão e limitações

Existem mudanças e limitações específicas em cada versão do PySpark. O Catalyst foi introduzido na 1.3 e melhorado significativamente na 2.0. Alguns métodos só chegaram em versões posteriores, como groupingSets().

Aqui está uma tabela com algumas das maiores mudanças funcionais no PySpark. Garanta que você usa a versão certa do Spark e do PySpark para os recursos de que precisa!

Recurso

Versão do Spark

Descrição

cube() / rollup()

1.4+

Agregações hierárquicas úteis em analytics estilo OLAP (3.4+ suporta Spark Connect)

pandas_udf

2.3+

UDFs vetorizadas usando Apache Arrow para execução mais rápida (3.4+ suporta Spark Connect, 4.0+ suporta SCALAR)

Adaptive Query Execution (AQE)

3.0+

Ajusta dinamicamente joins, shuffles e tratamento de skew em tempo de execução

Modo de compatibilidade ANSI SQL

3.0+

Relatos de erro e comportamento de expressões mais precisos

groupingSets()

4.0+

Permite múltiplos agrupamentos em uma única agregação (por exemplo, subtotais)

Conclusão

A função groupBy do PySpark é uma ferramenta essencial para agregação de dados em ambientes distribuídos. Seja para resumir dados por região, calcular métricas médias ou realizar análises complexas em múltiplos níveis, o groupBy oferece uma API escalável e flexível para workloads de big data.

Só não esqueça destas boas práticas para reduzir tempo de execução e custos:

  • Use funções nativas sempre que possível
  • Evite armadilhas de performance como shuffle excessivo e skew de dados.
  • Aproveite as otimizações e ferramentas de profiling do Spark para ajustar jobs grandes.

Conforme o PySpark evolui, podemos esperar otimizações ainda mais inteligentes, suporte nativo a agregações complexas e melhor integração com Pandas e sintaxe ao estilo SQL. É possível que traga integração ainda melhor com modelos em larga escala. Se você quer se aprofundar em PySpark, confira estes recursos da DataCamp:

Perguntas frequentes sobre PySpark groupBy

O que o método groupBy() do PySpark realmente retorna?

Ele retorna um objeto GroupedData, não um DataFrame. Esse objeto precisa ser seguido por um método de agregação como .count(), .sum() ou .agg() para retornar um novo DataFrame com os resultados do agrupamento.

O que é data skew e como ele afeta o groupBy()?

Data skew ocorre quando uma ou poucas chaves concentram dados demais, sobrecarregando alguns executors enquanto outros ficam ociosos. Isso pode levar a baixa performance ou falhas. Você pode mitigar o skew usando salting, particionamento customizado ou broadcast joins para dimensões pequenas.

Posso usar groupBy() em várias colunas no PySpark?

Sim, você pode agrupar por uma lista de colunas.

Quando devo usar funções de agregação personalizadas (UDFs ou pandas_udfs)?

Use-as apenas quando as funções nativas não derem conta. Funções nativas são mais rápidas e se beneficiam das otimizações do Spark. UDFs impedem que o Catalyst otimize o plano, enquanto pandas_udf troca escalabilidade por flexibilidade e deve ser usada só para datasets de pequeno a médio porte.

Como otimizo operações de groupBy() para performance no PySpark?

Minimize shuffles usando .repartition() e ajustando spark.sql.shuffle.partitions. Lembre-se de usar .cache() ao reutilizar resultados agregados. Verifique como o Spark está otimizando e suas recomendações com .explain().


Tim Lu's photo
Author
Tim Lu
LinkedIn

Sou um cientista de dados com experiência em análise espacial, machine learning e pipelines de dados. Trabalhei com GCP, Hadoop, Hive, Snowflake, Airflow e outros processos de engenharia/ciência de dados.

Tópicos
PySpark
Python

Principais cursos de PySpark

Curso

Fundamentos do PySpark

4 h
157.8K
Aprenda a implementar o gerenciamento de dados distribuídos e o machine learning no Spark usando o pacote PySpark.
Ver detalhesRight Arrow
Iniciar Curso
Ver maisRight Arrow
Relacionado

blog

Uma introdução aos polares: Ferramenta Python para análise de dados em grande escala

Explore o Polars, uma biblioteca Python robusta para manipulação e análise de dados de alto desempenho. Saiba mais sobre seus recursos, suas vantagens em relação ao pandas e como ele pode revolucionar seus processos de análise de dados.
Moez Ali's photo

Moez Ali

9 min

Tutorial

Tutorial do Pyspark: Primeiros passos com o Pyspark

Descubra o que é o Pyspark e como ele pode ser usado, com exemplos.
Natassha Selvaraj's photo

Natassha Selvaraj

10 min

Tutorial

Como usar GROUP BY e HAVING no SQL

Um guia intuitivo para você descobrir os dois comandos SQL mais populares para agregar linhas do seu conjunto de dados
Eugenia Anello's photo

Eugenia Anello

6 min

data-frames-in-python-banner_cgzjxy.jpeg

Tutorial

Pandas Tutorial: DataFrames em Python

Explore a análise de dados com Python. Os DataFrames do Pandas facilitam a manipulação de seus dados, desde a seleção ou substituição de colunas e índices até a remodelagem dos dados.
Karlijn Willems's photo

Karlijn Willems

15 min

Tutorial

Tutorial de seleção de colunas em Python

Use o Python Pandas e selecione colunas de DataFrames. Siga nosso tutorial com exemplos de código e aprenda diferentes maneiras de selecionar seus dados hoje mesmo!
DataCamp Team's photo

DataCamp Team

7 min

Clustering k-means

Tutorial

Introdução ao k-Means Clustering com o scikit-learn em Python

Neste tutorial, saiba como aplicar o k-Means Clustering com o scikit-learn em Python

Kevin Babitz

8 min

Ver MaisVer Mais