DEV Community

Cover image for Como aprendi Apache Spark — Parte 1: Por que o Spark mudou o processamento de Big Data
renanpyd
renanpyd

Posted on

Como aprendi Apache Spark — Parte 1: Por que o Spark mudou o processamento de Big Data

Hoje falamos naturalmente sobre Data Lakes, Lakehouses, processamento distribuído em cloud, Delta Lake, Apache Iceberg, pipelines em tempo real e plataformas como Databricks.

Naquele período, porém, uma palavra dominava grande parte das discussões sobre processamento distribuído:

Hadoop.

Mais especificamente, Hadoop MapReduce.

Entender esse contexto é importante porque Spark não surgiu simplesmente como "mais um framework de Big Data".

Ele nasceu tentando resolver limitações reais do processamento distribuído da época.

E foi justamente isso que chamou minha atenção.


Antes do Spark: o mundo era muito mais MapReduce

Imagine que temos alguns terabytes de logs e precisamos descobrir quantas vezes cada palavra aparece nesses dados.

Conceitualmente, podemos dividir o problema em duas etapas:

MAP
documento -> palavras -> (palavra, 1)

REDUCE
(palavra, 1), (palavra, 1), (palavra, 1)
                     ↓
              (palavra, 3)
Enter fullscreen mode Exit fullscreen mode

Essa ideia é extremamente poderosa.

Podemos dividir um dataset enorme entre diversas máquinas, processar partes dele paralelamente e depois combinar os resultados.

O Hadoop transformou esse modelo em uma plataforma capaz de processar volumes enormes de dados utilizando clusters de máquinas.

Mas existia um problema.

Para determinados workloads, principalmente aqueles compostos por múltiplas etapas, havia muita leitura e escrita de dados entre operações.

Imagine um algoritmo com várias etapas:

Dados
  ↓
Job 1
  ↓
Disco
  ↓
Job 2
  ↓
Disco
  ↓
Job 3
  ↓
Resultado
Enter fullscreen mode Exit fullscreen mode

Quando comecei a estudar processamento distribuído, uma das primeiras coisas que percebi foi:

Em Big Data, mover dados frequentemente custa mais do que executar a própria transformação.

Rede custa.

Disco custa.

Serialização custa.

Shuffle custa.

E essa observação continua extremamente relevante na Engenharia de Dados atual.


Então apareceu o Spark

Apache Spark nasceu como um projeto de pesquisa no AMPLab da University of California, Berkeley.

Uma das ideias centrais era permitir que determinados conjuntos de dados fossem reutilizados eficientemente durante diferentes etapas de uma computação distribuída.

Isso era especialmente interessante para workloads iterativos.

Pense, por exemplo, em algoritmos que repetem operações sobre os mesmos dados:

Dataset
   ↓
Iteração 1
   ↓
Iteração 2
   ↓
Iteração 3
   ↓
Iteração 4
   ↓
Resultado
Enter fullscreen mode Exit fullscreen mode

Machine Learning é um ótimo exemplo.

Muitos algoritmos precisam executar várias iterações sobre um mesmo conjunto de dados.

Se cada etapa precisar reconstruir todo o estado intermediário utilizando armazenamento persistente, o custo pode crescer rapidamente.

Spark propôs uma abstração muito interessante para esse problema.

Os RDDs.


RDD: Resilient Distributed Dataset

RDD significa:

Resilient Distributed Dataset

Em português, podemos pensar em algo como:

conjunto de dados distribuído e resiliente.

Mas essa tradução não explica por que a ideia foi tão importante.

Um RDD representa uma coleção de elementos distribuída pelas máquinas de um cluster e capaz de ser processada paralelamente.

Visualmente:

                    RDD
                     │
          ┌──────────┼──────────┐
          │          │          │
          ▼          ▼          ▼
      Partição 1  Partição 2  Partição 3
          │          │          │
          ▼          ▼          ▼
      Executor A  Executor B  Executor C
Enter fullscreen mode Exit fullscreen mode

Isso permite que diferentes partes do conjunto de dados sejam processadas simultaneamente.

Mas a parte mais interessante está no Resilient.


Por que "Resilient"?

Clusters falham.

Máquinas podem parar.

Executors podem morrer.

Processos podem ser encerrados.

Nós podem ficar indisponíveis.

Uma plataforma de processamento distribuído precisa assumir que falhas vão acontecer.

Spark consegue reconstruir partições perdidas utilizando informações sobre como aquele dataset foi produzido.

Essa cadeia de dependências é chamada de:

lineage.

Imagine:

Arquivo
   │
   ▼
textFile()
   │
   ▼
  RDD A
   │
   │ filter()
   ▼
  RDD B
   │
   │ map()
   ▼
  RDD C
Enter fullscreen mode Exit fullscreen mode

Spark conhece essa sequência de transformações.

Se uma partição de RDD C for perdida, ele não necessariamente precisa manter uma cópia completa de tudo.

Em muitos casos, pode reconstruir a partição a partir do lineage.

Essa ideia é fundamental para entender a arquitetura do Spark.


Um primeiro RDD em PySpark

Vamos olhar um exemplo extremamente simples.

numeros = sc.parallelize([1, 2, 3, 4, 5])
Enter fullscreen mode Exit fullscreen mode

Aqui criamos um RDD contendo cinco elementos.

Agora podemos aplicar uma transformação:

dobrados = numeros.map(lambda x: x * 2)
Enter fullscreen mode Exit fullscreen mode

Intuitivamente, poderíamos imaginar:

[1, 2, 3, 4, 5]

       map(x * 2)

[2, 4, 6, 8, 10]
Enter fullscreen mode Exit fullscreen mode

Mas existe um detalhe fundamental.

Quando executamos:

dobrados = numeros.map(lambda x: x * 2)
Enter fullscreen mode Exit fullscreen mode

Spark não precisa calcular imediatamente o resultado.

Ele registra a transformação.

Esse comportamento é chamado de:

Lazy Evaluation

E esse é um dos conceitos mais importantes de toda a arquitetura Spark.


Spark não executa tudo imediatamente

Considere:

numeros = sc.parallelize([1, 2, 3, 4, 5])

pares = numeros.filter(lambda x: x % 2 == 0)

dobrados = pares.map(lambda x: x * 2)
Enter fullscreen mode Exit fullscreen mode

Podemos imaginar o plano:

parallelize
     │
     ▼
   filter
     │
     ▼
    map
Enter fullscreen mode Exit fullscreen mode

Até aqui estamos descrevendo o que queremos fazer.

Agora executamos:

resultado = dobrados.collect()

print(resultado)
Enter fullscreen mode Exit fullscreen mode

Resultado:

[4, 8]
Enter fullscreen mode Exit fullscreen mode

collect() é uma ação.

É nesse momento que Spark precisa efetivamente produzir o resultado.

Essa separação nos leva a dois conceitos fundamentais.


Transformations vs Actions

Grande parte da programação com Spark pode ser compreendida através dessa divisão.

Transformations

Transformations criam um novo dataset a partir de outro.

Exemplos:

rdd.map(...)
rdd.filter(...)
rdd.flatMap(...)
rdd.distinct()
Enter fullscreen mode Exit fullscreen mode

Normalmente elas são avaliadas de forma lazy.

Podemos imaginar:

RDD
 │
 ├── filter()
 │
 ▼
RDD
 │
 ├── map()
 │
 ▼
RDD
Enter fullscreen mode Exit fullscreen mode

Estamos construindo um plano de transformação.


Actions

Actions solicitam algum resultado da computação.

Exemplos clássicos:

rdd.count()
rdd.collect()
rdd.first()
rdd.take(10)
Enter fullscreen mode Exit fullscreen mode

Quando uma Action é executada, Spark precisa avaliar as transformações necessárias para produzir o resultado.

Por exemplo:

resultado = (
    sc.parallelize(range(100))
      .filter(lambda x: x % 2 == 0)
      .map(lambda x: x * 10)
      .take(5)
)
Enter fullscreen mode Exit fullscreen mode

Temos:

Dados
  │
  ▼
filter()
  │
  ▼
map()
  │
  ▼
take(5)
  │
  ▼
Resultado
Enter fullscreen mode Exit fullscreen mode

Esse modelo permite ao Spark entender uma sequência de operações antes da execução.

Mais tarde veremos como isso se relaciona com jobs, stages, tasks e DAGs.


O clássico WordCount

Nenhuma viagem pela história do processamento distribuído estaria completa sem ele.

O famoso:

WordCount.

linhas = sc.textFile("dados.txt")

palavras = linhas.flatMap(
    lambda linha: linha.split()
)

pares = palavras.map(
    lambda palavra: (palavra, 1)
)

contagem = pares.reduceByKey(
    lambda a, b: a + b
)

resultado = contagem.collect()
Enter fullscreen mode Exit fullscreen mode

Conceitualmente:

Arquivo
   │
   ▼
Linhas
   │
   │ flatMap()
   ▼
Palavras
   │
   │ map()
   ▼
(palavra, 1)
   │
   │ reduceByKey()
   ▼
(palavra, quantidade)
Enter fullscreen mode Exit fullscreen mode

Esse pequeno exemplo contém vários conceitos importantes:

  • leitura distribuída;
  • transformação dos dados;
  • criação de pares chave/valor;
  • agregação;
  • execução distribuída;
  • movimentação de dados entre partições.

E existe uma palavra especialmente importante nessa lista:

movimentação.

Porque reduceByKey() pode exigir que registros com a mesma chave sejam reunidos.

Isso nos leva a um dos assuntos mais importantes de performance em Spark:

shuffle.

Vamos dedicar um artigo específico a isso mais adiante.


Spark não é simplesmente "Hadoop mais rápido"

Durante muito tempo tornou-se comum encontrar explicações como:

"Spark é Hadoop em memória."

Essa descrição é conveniente.

Mas é incompleta.

Spark não é apenas uma implementação mais rápida de MapReduce.

Ele oferece um modelo de execução diferente e uma API capaz de representar pipelines de transformação muito mais naturalmente.

Compare conceitualmente:

MapReduce

Job
 ↓
Disco
 ↓
Job
 ↓
Disco
 ↓
Job
Enter fullscreen mode Exit fullscreen mode

com:

Spark

Transformação
     ↓
Transformação
     ↓
Transformação
     ↓
Action
Enter fullscreen mode Exit fullscreen mode

Essa mudança no modelo de programação é tão importante quanto a questão de performance.


Spark também não mantém "tudo em memória"

Outro erro que encontrei muitas vezes ao aprender Spark foi:

"Spark coloca todos os dados na RAM."

Não.

Spark pode manter dados em memória quando isso é útil e possível, mas isso não significa que todo processamento acontece exclusivamente na memória.

Dependendo da operação, configuração e volume de dados, Spark pode utilizar:

  • memória;
  • disco;
  • rede;
  • armazenamento externo;
  • recomputação.

O modelo real é muito mais sofisticado.

E entender isso evita várias decisões ruins de arquitetura.


O que continua relevante em 2026?

Hoje, na maior parte dos pipelines modernos, provavelmente não começaremos escrevendo diretamente RDDs.

É muito mais comum trabalharmos com:

df = spark.read.parquet("/dados/clientes")
Enter fullscreen mode Exit fullscreen mode

e depois:

resultado = (
    df
    .filter(df.ativo == True)
    .groupBy("estado")
    .count()
)
Enter fullscreen mode Exit fullscreen mode

Ou utilizando SQL:

SELECT
    estado,
    COUNT(*) AS quantidade
FROM clientes
WHERE ativo = true
GROUP BY estado;
Enter fullscreen mode Exit fullscreen mode

DataFrames e Spark SQL permitem que o Spark realize otimizações muito mais sofisticadas sobre as consultas.

Então por que estudar RDD?

Porque vários conceitos fundamentais continuam aparecendo por baixo das abstrações modernas:

DataFrame / SQL
       │
       ▼
Plano lógico
       │
       ▼
Otimização
       │
       ▼
Plano físico
       │
       ▼
Stages
       │
       ▼
Tasks
       │
       ▼
Executors
Enter fullscreen mode Exit fullscreen mode

Quando alguma coisa fica lenta, cara ou instável, conhecer apenas:

df.groupBy(...)
Enter fullscreen mode Exit fullscreen mode

geralmente não é suficiente.

Precisamos entender o que está acontecendo por baixo da API.


Uma lição que continua válida

Uma das coisas mais importantes que aprendi estudando Spark foi que Engenharia de Dados distribuída não consiste apenas em escrever transformações.

Precisamos pensar em:

Dados
  ↓
Particionamento
  ↓
Transformações
  ↓
Dependências
  ↓
Shuffle
  ↓
Stages
  ↓
Tasks
  ↓
Executors
  ↓
CPU / Memória / Disco / Rede
Enter fullscreen mode Exit fullscreen mode

Quando começamos a enxergar um pipeline dessa forma, vários problemas deixam de parecer misteriosos.

Um job lento deixa de ser simplesmente:

"Spark está lento."

E passa a gerar perguntas melhores:

Houve shuffle?

Existe data skew?

As partições estão bem dimensionadas?

Estamos movimentando dados desnecessariamente?

Algum executor está sofrendo pressão de memória?

O driver está recebendo dados demais?

Estamos usando a abstração correta?

Essas perguntas continuam extremamente atuais.


Daqui para frente

Neste artigo vimos os fundamentos que me fizeram começar a olhar Spark de forma diferente:

  • o contexto do Hadoop MapReduce;
  • processamento distribuído;
  • RDDs;
  • resiliência;
  • lineage;
  • transformations;
  • actions;
  • lazy evaluation;
  • e a ideia de distribuir uma computação entre várias máquinas.

Agora podemos entrar de verdade na arquitetura.

No próximo artigo vamos responder uma pergunta fundamental:

Quando executo um programa Spark, quem realmente faz o trabalho?

Vamos conhecer:

Driver, SparkContext, Cluster Manager, Executors, Jobs, Stages e Tasks.

E vamos montar o caminho completo:

Aplicação
    ↓
Driver
    ↓
Cluster Manager
    ↓
Executors
    ↓
Stages
    ↓
Tasks
    ↓
Dados
Enter fullscreen mode Exit fullscreen mode

Entender esse fluxo muda completamente a maneira como diagnosticamos performance, memória e falhas em aplicações Spark.


Série: Como aprendi Apache Spark

← Parte 0: Como aprendi Apache Spark: revisitando uma jornada pela Engenharia de Dados

Próximo: Parte 2 — Por dentro da arquitetura do Apache Spark: Driver, Executors, Jobs, Stages e Tasks

Top comments (0)