Transformações no PySpark: o que você precisa saber
Transformações no PySpark
Transformações no PySpark: o que você precisa saber
Transformações no PySpark
Nesse artigo falarei um pouco sobre as transformations (transformações) no PySpark e irei aprofundar o assunto nos próximos artigos da série.
Estrutura de dados imutáveis e transformations
No Spark, quando criamos um DataFrame ou qualquer estrutura de dados, ela não pode ser modificada — isso significa que são imutáveis. Mas aí vem a pergunta: como alteramos um DataFrame? A resposta é simples: você não altera, simplesmente instrui o Spark a criar um novo DataFrame a partir do anterior, descrevendo as transformações necessárias. Essas instruções são chamadas de transformations (transformações).
Exemplo:
from pyspark.sql.functions import col
brazil_flights = flights_df.filter(col("DEST_COUNTRY_NAME") == "Brazil")
Perceba que no exemplo acima não alteramos o DataFrame flights_df, apenas criamos um DataFrame novo chamado brazil_flights aplicando um filtro.
Na prática ocorreu o seguinte:
DataFrame original → [transformação] → Novo DataFrame
(imutável) (também imutável)
↓
Não foi alterado!
Lazy Evaluation
Um detalhe importante é que o Spark não executa nada nesse momento, e isso se deve ao Lazy Evaluation, ou seja, o Spark não executa nada até que uma action (falaremos sobre isso em outro momento) seja executada.
Com lazy evaluation, o Spark recebe todas as transformações primeiro, monta o plano completo, e só então analisa o que realmente precisa ser feito. É nessa análise que ele percebe: “espera, no final desse pipeline só precisamos de uma linha — então nem faz sentido ler o resto dos dados”. E empurra esse filtro para o começo.
Em resumo: você nunca executa transformações diretamente — você descreve o que quer fazer, e o Spark decide a melhor forma de executar tudo de uma vez.
Transformações mais comuns do dia a dia
Agora que entendemos o conceito de transformações e lazy evaluation, vamos ver na prática as operações que você vai usar com mais frequência no seu dia a dia com Spark. Usaremos um DataFrame de exemplo com dados de voos entre países.
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("transformacoes").getOrCreate()
df = spark.read.format("json").load("/data/flight-data/json/2015-summary.json")
df.show(3)
+-----------------+-------------------+-----+
|DEST_COUNTRY_NAME|ORIGIN_COUNTRY_NAME|count|
+-----------------+-------------------+-----+
| United States| Romania| 15|
| United States| Croatia| 1|
| United States| Ireland| 344|
+-----------------+-------------------+-----+
1. select — Selecionar colunas
O select é provavelmente a transformação que você mais vai usar. Funciona exatamente como o SELECT do SQL: você escolhe quais colunas quer trabalhar.
df.select("DEST_COUNTRY_NAME", "ORIGIN_COUNTRY_NAME").show(3)
+-----------------+-------------------+
|DEST_COUNTRY_NAME|ORIGIN_COUNTRY_NAME|
+-----------------+-------------------+
| United States| Romania|
| United States| Croatia|
| United States| Ireland|
+-----------------+-------------------+
Com o selectExpr, você pode ir além e escrever expressões SQL diretamente como string, o que é muito útil para renomear colunas ou criar novas na hora:
df.selectExpr(
"DEST_COUNTRY_NAME as destino",
"ORIGIN_COUNTRY_NAME as origem",
"(DEST_COUNTRY_NAME = ORIGIN_COUNTRY_NAME) as voo_domestico"
).show(3)
+-------------+---------+-------------+
| destino| origem|voo_domestico|
+-------------+---------+-------------+
|United States| Romania| false|
|United States| Croatia| false|
|United States| Ireland| false|
+-------------+---------+-------------+
2. filter / where — Filtrar linhas
Para filtrar linhas, usamos filter ou where — os dois fazem a mesma coisa. Por ser mais familiar para quem vem do SQL, o where costuma ser preferido.
# Filtrando voos com contagem menor que 2
df.where("count < 2").show(3)
+-----------------+-------------------+-----+
|DEST_COUNTRY_NAME|ORIGIN_COUNTRY_NAME|count|
+-----------------+-------------------+-----+
| United States| Croatia| 1|
| United States| Singapore| 1|
| Moldova| United States| 1|
+-----------------+-------------------+-----+
Você também pode encadear múltiplos filtros. O Spark é inteligente o suficiente para juntá-los em uma única operação internamente:
from pyspark.sql.functions import col
df.where(col("count") < 2)\
.where(col("ORIGIN_COUNTRY_NAME") != "Croatia")\
.show(3)
Dica: Encadear
whereé mais legível do que escrever um único filtro complexo, e o Spark aplica os dois ao mesmo tempo — sem penalidade de performance.
3. withColumn — Adicionar ou transformar colunas
O withColumn serve tanto para criar uma nova coluna quanto para modificar uma existente. Ele recebe dois argumentos: o nome da coluna e a expressão que define seu valor.
from pyspark.sql.functions import expr
# Criando uma coluna que indica se o voo é doméstico
df.withColumn(
"voo_domestico",
expr("ORIGIN_COUNTRY_NAME = DEST_COUNTRY_NAME")
).show(3)
+-----------------+-------------------+-----+-------------+
|DEST_COUNTRY_NAME|ORIGIN_COUNTRY_NAME|count|voo_domestico|
+-----------------+-------------------+-----+-------------+
| United States| Romania| 15| false|
| United States| Croatia| 1| false|
| United States| Ireland| 344| false|
+-----------------+-------------------+-----+-------------+
Lembre-se: como as estruturas são imutáveis, o DataFrame original não foi alterado. O withColumn retorna um novo DataFrame com a coluna adicionada.
4. withColumnRenamed — Renomear colunas
Para renomear uma coluna, usamos o withColumnRenamed. O primeiro argumento é o nome atual e o segundo é o novo nome:
df.withColumnRenamed("DEST_COUNTRY_NAME", "destino")\
.withColumnRenamed("ORIGIN_COUNTRY_NAME", "origem")\
.show(3)
+-------------+---------+-----+
| destino| origem|count|
+-------------+---------+-----+
|United States| Romania| 15|
|United States| Croatia| 1|
|United States| Ireland| 344|
+-------------+---------+-----+
5. drop — Remover colunas
Às vezes é mais fácil remover colunas que não precisamos do que selecionar todas as que queremos. Para isso usamos o drop:
df.drop("ORIGIN_COUNTRY_NAME").show(3)
+-----------------+-----+
|DEST_COUNTRY_NAME|count|
+-----------------+-----+
| United States| 15|
| United States| 1|
| United States| 344|
+-----------------+-----+
Você também pode remover múltiplas colunas de uma vez:
df.drop("ORIGIN_COUNTRY_NAME", "DEST_COUNTRY_NAME").show(3)
6. orderBy / sort — Ordenar dados
Para ordenar um DataFrame, usamos orderBy (ou sort, que é equivalente). Por padrão a ordenação é crescente, mas podemos controlar isso com asc() e desc():
from pyspark.sql.functions import desc, asc
# Ordenando por count decrescente e depois por nome do destino crescente
df.orderBy(col("count").desc(), col("DEST_COUNTRY_NAME").asc()).show(5)
+-----------------+-------------------+------+
|DEST_COUNTRY_NAME|ORIGIN_COUNTRY_NAME| count|
+-----------------+-------------------+------+
| United States| United States|370002|
| United States| Canada| 8483 |
| Canada| United States| 8399 |
| United States| Mexico| 7187 |
| Mexico| United States| 7140 |
+-----------------+-------------------+------+
7. limit — Limitar resultados
O limit restringe quantas linhas o DataFrame retorna — muito útil em conjunto com o orderBy para pegar os top N registros:
# Top 3 rotas com mais voos
df.orderBy(col("count").desc()).limit(3).show()
+-----------------+-------------------+------+
|DEST_COUNTRY_NAME|ORIGIN_COUNTRY_NAME| count|
+-----------------+-------------------+------+
| United States| United States|370002|
| United States| Canada| 8483|
| Canada| United States| 8399|
+-----------------+-------------------+------+
Encadeando transformações
Uma das grandes vantagens do Spark é poder encadear todas essas transformações de forma fluente. Graças ao lazy evaluation, o Spark só executa tudo quando você chama uma action (como o show()):
resultado = df\
.where(col("ORIGIN_COUNTRY_NAME") != "United States")\
.withColumn("voo_domestico", expr("ORIGIN_COUNTRY_NAME = DEST_COUNTRY_NAME"))\
.withColumnRenamed("DEST_COUNTRY_NAME", "destino")\
.drop("ORIGIN_COUNTRY_NAME")\
.orderBy(col("count").desc())\
.limit(5)
resultado.show()
+-------------+-----+-------------+
| destino|count|voo_domestico|
+-------------+-----+-------------+
| Canada| 8399| false|
| Mexico| 7140| false|
|Great Britain| 2025| false|
| Japan| 1548| false|
| Germany| 1468| false|
+-------------+-----+-------------+
Todo esse pipeline — filtrar, criar coluna, renomear, remover, ordenar e limitar — é montado como um plano de execução e processado de uma só vez quando o show() é chamado.
No próximo artigo da série, vamos aprofundar o tema das transformações e entender a diferença entre Narrow vs. Wide Transformations — um conceito fundamental para compreender como o Spark distribui e processa os dados no cluster. Até lá!
메타데이터
- post_id
- 02e33f2acecb
- slug
- transformações-no-pyspark-o-que-você-precisa-saber-02e33f2acecb
- url
- https://medium.com/@heycaio/transforma%C3%A7%C3%B5es-no-pyspark-o-que-voc%C3%AA-precisa-saber-02e33f2acecb
- canonical_url
- https://medium.com/@heycaio/transforma%C3%A7%C3%B5es-no-pyspark-o-que-voc%C3%AA-precisa-saber-02e33f2acecb
- author_url
- https://medium.com/@heycaio
- status
- ok
- fetched_at
- 2026-06-09 15:37:30