Gustavo Ramos
Voltar para o blog
Aprofundado

Construindo um Pipeline de ETL em Python Pronto para Produção (Sem Superengenharia)

Publicado em25 de junho de 20265 min de leitura
pythonetlpandaspostgresqldata-pipelines

Um script que busca dados de uma API, limpa e escreve num banco é fácil de escrever e fácil de errar de formas que só aparecem quando ele já está rodando sozinho às três da manhã. A distância entre "funciona na minha máquina, uma vez" e "roda toda noite sem ninguém olhando" não é mais código — é um punhado de decisões específicas tomadas com antecedência sobre falha, não sobre o caminho feliz.

Extract: trate a API como não confiável por padrão

Todo pipeline que consome uma API de terceiro eventualmente vai bater em um timeout, um rate limit ou um formato de resposta sutilmente diferente do que a documentação prometia. A etapa de extração deve assumir isso desde o primeiro dia, em vez de tratamento de erro colado depois da primeira queda.

import time
import requests

def fetch_page(url: str, params: dict, max_retries: int = 3) -> dict:
    for attempt in range(1, max_retries + 1):
        response = requests.get(url, params=params, timeout=10)
        if response.status_code == 429:
            wait = int(response.headers.get("Retry-After", 2 ** attempt))
            time.sleep(wait)
            continue
        response.raise_for_status()
        return response.json()
    raise RuntimeError(f"Failed to fetch {url} after {max_retries} attempts")

Os dois detalhes que mais importam aqui são o timeout explícito (uma requisição travada sem ele pode travar o pipeline inteiro indefinidamente) e tratar 429 como um caso distinto de uma falha genérica — a maioria das APIs te diz exatamente quanto tempo esperar, e ignorar esse header é como você transforma um rate limit em um banimento.

Transform: valide na borda, não no meio da sua lógica

O Pandas te tenta a escrever lógica de transformação que assume entrada limpa, porque na maior parte do tempo a entrada é limpa — até um registro ter um nulo onde você não esperava, e um KeyError três funções mais adiante matar a execução inteira em silêncio. Empurre a validação para a borda, logo após a extração, para que o resto do pipeline confie nos próprios dados:

import pandas as pd

REQUIRED_COLUMNS = {"id", "created_at", "amount"}

def load_and_validate(records: list[dict]) -> pd.DataFrame:
    df = pd.DataFrame(records)
    missing = REQUIRED_COLUMNS - set(df.columns)
    if missing:
        raise ValueError(f"API response missing required columns: {missing}")

    df["created_at"] = pd.to_datetime(df["created_at"], errors="coerce")
    bad_dates = df["created_at"].isna().sum()
    if bad_dates:
        df = df.dropna(subset=["created_at"])

    return df

Se você descarta linhas ruins, coloca em quarentena ou falha a execução inteira depende do domínio — descartar em silêncio é aceitável para um agregado analítico noturno, inaceitável para registros financeiros. O ponto não é a política específica; é tornar essa decisão explícita e visível no código, em vez de deixar NaNs se propagarem silenciosamente para agregações mais adiante.

Load: torne as reexecuções seguras

A propriedade mais valiosa que um pipeline pode ter é idempotência — rodá-lo duas vezes com a mesma entrada não deve gerar linhas duplicadas nem estado corrompido. Um INSERT sem estratégia de conflito falha nisso na primeira vez que um job é reexecutado após uma falha parcial:

from sqlalchemy import text

def upsert_batch(engine, df: pd.DataFrame, table: str) -> None:
    with engine.begin() as conn:
        for row in df.to_dict(orient="records"):
            conn.execute(
                text(f"""
                    INSERT INTO {table} (id, created_at, amount)
                    VALUES (:id, :created_at, :amount)
                    ON CONFLICT (id) DO UPDATE
                    SET created_at = EXCLUDED.created_at,
                        amount = EXCLUDED.amount
                """),
                row,
            )

ON CONFLICT ... DO UPDATE (o upsert do Postgres) significa que uma reexecução depois de uma instabilidade de rede ou uma falha parcial de lote simplesmente reaplica as mesmas linhas em vez de duplicá-las. Combinado com envolver o lote numa transação (engine.begin()), uma falha no meio do lote reverte de forma limpa em vez de deixar a tabela pela metade.

Agendamento: logue o suficiente para depurar uma falha que você não vai ver acontecer

O modo de falha que realmente importa para um pipeline sem supervisão não é "ele quebrou" — o cron ou o seu agendador vão te avisar disso. É "ele processou zero linhas em silêncio" ou "ele processou as linhas erradas em silêncio", que parecem idênticos a um sucesso num log que só diz Done.. Logue contagens de linhas em cada etapa, não só timestamps de início/fim, para que uma execução ruim seja diagnosticável só pelo log, sem precisar reproduzir o problema localmente na manhã seguinte.

Nada disso é exótico. É a diferença entre um script e um pipeline: o script assume que o mundo coopera, o pipeline assume que não vai cooperar, e torna essa suposição barata de tratar em vez de cara de descobrir.

Gerado por Claude Sonnet 5 (seeded manually via Claude Code) · 25 de junho de 2026 · verificado em build antes da publicação