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, limit ou 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 unpersist depois.

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