HenriqueDePaula12
data_engineer
Jupyter Notebook

Esse é o meu primeiro repo tratando de fim a fim, uma pipeline de dados abertos do governo brasileiro relacionado a compras de contrato e cronogramas anuais com spark, em pyspark e SQL!

Last updated Apr 4, 2022
10
Stars
1
Forks
0
Issues
0
Stars/day
Attention Score
0
Language breakdown
Jupyter Notebook 99.7%
Python 0.2%
Shell 0.0%
Files click to expand
README

Olá!

Esse é o meu primeiro repo tratando de fim a fim, uma pipeline de dados abertos do governo brasileiro relacionado a compras de contrato e cronogramas anuais com spark, em pyspark e SQL!

O código se encontra aqui e o dado pode ser obtido por meio desse link

from pyspark.sql import SparkSession

##################################################### VARIABLES #####################################################

PATHLANDINGZONE_CSV = '../datalake/landing/comprasnet-contratos-anual-cronogramas-latest.csv' PATHPROCESSINGZONE = '../datalake/processing' PATHCURATEDZONE = '../datalake/curated'

##################################################### QUERY #########################################################

QUERY = """

WITH tmp as ( SELECT cast(id as integer) as id, cast(contratoid as integer) as contratoid, tipo, numero, receita_despesa, observacao, mesref, anoref, cast(vencimento as date) as vencimento, retroativo, cast(valor as decimal (10,2)) as valor, year(vencimento) as year, month(vencimento) as month, dayofmonth(vencimento) as day FROM df ) SELECT * FROM tmp WHERE year = 2021 OR year = 2022 ORDER BY year desc

"""

##################################################### SCRIPT #########################################################

def csvtoparquet(spark, pathcsv, pathparquet): df = spark.read.option('header', True).csv(path_csv) return df.write.mode('overwrite').format('parquet').save(path_parquet)

def createview(spark, pathparquet): df = spark.read.parquet(path_parquet) df.createOrReplaceTempView('df')

def writecurated(spark, pathcurated): df2 = spark.sql(QUERY) ( df2 .orderBy('year', ascending=False) .orderBy('month', ascending=False) .orderBy('day', ascending=False) .write.partitionBy('year','month','day') .mode('overwrite') .format('parquet') .save(path_curated) )

if name == "main": spark = ( SparkSession.builder .master("local[*]") .getOrCreate() )

spark.sparkContext.setLogLevel("ERROR") csvtoparquet(spark, PATHLANDINGZONECSV, PATHPROCESSING_ZONE)

createview(spark, PATHPROCESSING_ZONE) writecurated(spark, PATHCURATED_ZONE )

  • Basicamente, extraimos os dados para a zona landing, depois, escrevemos o mesmo dado em diferente formato na zona processing, no caso parquet, por se tratar de um formato otimizado e mais leve.
  • Após, criamos uma view do dado recém salvo na zona processing, já em parquet, que otimiza a leitura do spark, aplicamos uma query de transformação que enriquece o schema do dado e seleciona apenas os dados de 2021 e 2022, já pronto para ser consumido.
  • E por fim, escrevemos na zona curated o dado já tratado, enriquecido, particionado por ano, mês e dia e pronto para consumo.
Para rodar o script, basicamente você pode fazer no terminal:
spark-submit etl.py

Você também encontrará o mesmo código e ideia de ETL em notebooks, em versão pyspark ou spark-sql.

Espero que gostem!

Qualquer dúvida, entrar em contato pelo LinkedIn.

:)

🔗 More in this category

© 2026 GitRepoTrend · HenriqueDePaula12/data_engineer · Updated daily from GitHub