Processamento distribuído
Tutorial PySpark: do DataFrame ao job que aguenta terabytes
PySpark é a habilidade que separa quem processa planilhas de quem processa a empresa. O caminho abaixo é o mesmo que você usa em Databricks, EMR ou Dataproc — só o tamanho do cluster muda.
Começando: SparkSession e leitura
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = (
SparkSession.builder
.appName("vendas")
.master("local[*]") # em produção o cluster define isso
.getOrCreate()
)
vendas = (
spark.read
.option("header", True)
.option("inferSchema", True)
.csv("data/vendas.csv")
)
vendas.printSchema()
vendas.show(5, truncate=False)Em produção, evite inferSchema: ele lê o arquivo duas vezes. Declare o schema explicitamente e ganhe tempo e previsibilidade.
As transformações do dia a dia
pagas = (
vendas
.filter(F.col("status") == "paga")
.withColumn("receita", F.col("quantidade") * F.col("preco_unitario"))
.withColumn("mes", F.date_trunc("month", F.col("data_pedido")))
)
resumo = (
pagas
.groupBy("mes", "categoria")
.agg(
F.sum("receita").alias("receita_total"),
F.countDistinct("cliente_id").alias("clientes"),
F.avg("receita").alias("ticket_medio"),
)
.orderBy(F.col("mes").desc(), F.col("receita_total").desc())
)
resumo.show()Nada disso rodou ainda. Só no show() o Spark executa — e executa o plano já otimizado, empurrando o filtro para o mais perto possível da leitura.
Joins sem estourar o cluster
clientes = spark.read.parquet("data/clientes")
# tabela pequena? force broadcast e evite shuffle
enriquecido = pagas.join(
F.broadcast(clientes.select("cliente_id", "uf", "segmento")),
on="cliente_id",
how="left",
)
# checando chaves órfãs antes de confiar no resultado
orfaos = pagas.join(clientes, "cliente_id", "left_anti").count()
print("pedidos sem cliente:", orfaos)left_anti é o teste de qualidade mais barato que existe: mostra o que não casou, sem duplicar linhas.
Window functions: ranking e comparação temporal
from pyspark.sql.window import Window
w_cliente = Window.partitionBy("cliente_id").orderBy("data_pedido")
historico = (
pagas
.withColumn("nro_pedido", F.row_number().over(w_cliente))
.withColumn("pedido_anterior", F.lag("receita").over(w_cliente))
.withColumn(
"variacao",
F.round((F.col("receita") - F.col("pedido_anterior")) / F.col("pedido_anterior") * 100, 1),
)
.withColumn(
"receita_acumulada",
F.sum("receita").over(w_cliente.rowsBetween(Window.unboundedPreceding, 0)),
)
)Escrita particionada e performance
(
resumo
.repartition("mes")
.write
.mode("overwrite")
.partitionBy("mes")
.parquet("s3://silver/vendas_resumo")
)- Small files: milhares de arquivos minúsculos matam a leitura. Ajuste partições antes de escrever.
- Skew: uma chave com 80% dos dados trava um executor. Salgue a chave ou use broadcast.
- collect(): traz tudo para o driver e derruba o job. Use
show,limitou escreva o resultado. - UDF Python: lenta por natureza. Prefira funções nativas de
pyspark.sql.functions. - cache(): só quando o mesmo DataFrame é reutilizado várias vezes — e sempre com
unpersistdepois.
Perguntas frequentes
- O que é PySpark?
- É a API Python do Apache Spark. Você escreve código Python e o Spark distribui o processamento por várias máquinas, permitindo tratar volumes que não cabem na memória de um notebook.
- PySpark é diferente de pandas?
- Sim. pandas processa em uma máquina e executa na hora; PySpark é distribuído e lazy — nada roda até você chamar uma ação como show, count ou write. Use pandas até alguns GB e PySpark acima disso.
- Preciso de um cluster para aprender PySpark?
- Não. Instale com pip install pyspark e rode em modo local usando todos os núcleos da sua máquina. A API é idêntica à do cluster; só o volume muda.
- O que é lazy evaluation no Spark?
- Transformações como select e filter só montam um plano de execução. Quando você chama uma ação, o Catalyst otimiza esse plano inteiro de uma vez — por isso a ordem em que você escreve importa menos do que se imagina.
- Qual a diferença entre repartition e coalesce?
- repartition embaralha os dados pela rede e pode aumentar ou reduzir partições; coalesce só reduz, sem shuffle completo. Use coalesce antes de escrever poucos arquivos e repartition quando precisar redistribuir dados desbalanceados.
- PySpark vale a pena com Databricks no mercado brasileiro?
- Vale, e muito: Databricks é dominante nas vagas brasileiras e roda Spark por baixo. Quem sabe PySpark aproveita quase tudo em Databricks, EMR, Dataproc ou Fabric.
Pronto para assinar?
7 dias grátis. Depois, menos que um café por mês — cancele quando quiser.
Assinar — 7 dias grátis- Sem cartão no teste grátis
- Cancele quando quiser
- 300+ exercícios
- 14 cursos completos