Como extrair mudanças incrementais em APIs de Dados: passo a passo
Aprende a extrair apenas os registos novos ou alterados de uma API de Dados — útil para reduzir tráfego, acelerar ETL e manter uma cópia local actualizada. Vamos explicar o porquê e mostrar um exemplo em Python que detecta e processa mudanças incrementais.
Pré-requisitos
- Python 3.8+ instalado
- Bibliotecas: requests e sqlite3 (integrada no Python)
- Uma API REST que exponha um campo de data/timestamp ou um campo incremental (ex.: updated_at ou id)
- Noções básicas de HTTP e JSON
Passo 1: Entender o modelo de mudança da API
Antes de escrever código, verifica se a API tem um campo que indica alterações (por exemplo updated_at) ou um identificador incremental. Se tiver updated_at em UTC, podes pedir apenas registos depois de uma data. Se só houver id incremental, podes pedir id > último_id. Saber isto evita puxar todos os dados e é a base da extracção incremental.
Passo 2: Estruturar um armazenamento de controlo (state)
Precisas de guardar o último timestamp ou último id extraído. Um ficheiro simples ou uma pequena base de dados SQLite funciona bem para um fluxo ETL simples. Aqui usamos SQLite para persistência segura.
import sqlite3
def init_state_db(path='state.db'):
conn = sqlite3.connect(path)
cur = conn.cursor()
cur.execute('''CREATE TABLE IF NOT EXISTS etl_state (
key TEXT PRIMARY KEY,
value TEXT
)''')
conn.commit()
return conn
def get_state(conn, key):
cur = conn.cursor()
cur.execute('SELECT value FROM etl_state WHERE key=?', (key,))
row = cur.fetchone()
return row[0] if row else None
def set_state(conn, key, value):
cur = conn.cursor()
cur.execute('REPLACE INTO etl_state (key, value) VALUES (?, ?)', (key, value))
conn.commit()
Passo 3: Fazer uma chamada condicionada à API
Constrói a query usando o valor guardado. Exemplo com parâmetro updated_after (muitas APIs aceitam algo semelhante). Trata erros comuns como 400/500 e limitações de rate limit; se faltar suporte a filtragem, terás de usar paginação e filtrar localmente.
import requests
from datetime import datetime
API_URL = 'https://api.exemplo.com/items'
def fetch_incremental(updated_after=None, page=1):
params = {'page': page}
if updated_after:
params['updated_after'] = updated_after # nome do parâmetro depende da API
resp = requests.get(API_URL, params=params, timeout=10)
resp.raise_for_status()
return resp.json() # assume JSON com lista e metadados
Passo 4: Processar resultados e actualizar o state
Ao receber os registos, processa-os (gravar na tua base local, transformar, etc.) e calcula o novo valor de state: max(updated_at) ou max(id). Só depois de gravar com sucesso actualiza o state para evitar perda de dados.
def process_and_update_state(conn, items):
# Exemplo mínimo: gravar no ficheiro ou base local (omiti implementação de gravação completa)
# Supõe que cada item tem 'id' e 'updated_at' em ISO 8601
if not items:
return None
max_ts = max(item['updated_at'] for item in items)
# aqui gravávamos os items numa tabela ou ficheiros
set_state(conn, 'last_updated_at', max_ts)
return max_ts
Passo 5: Juntar tudo num ciclo seguro
Executa um ciclo que lê o state, chama a API em páginas até não haver mais, processa e actualiza o state no fim. Trata erros temporários com retries simples e respeita rate limits com backoff linear ou exponencial.
import time
def run_once(conn):
last = get_state(conn, 'last_updated_at')
page = 1
all_items = []
while True:
try:
data = fetch_incremental(updated_after=last, page=page)
except requests.HTTPError as e:
print('Erro HTTP:', e)
break
items = data.get('items', [])
if not items:
break
# lógica de processamento aqui; acumulamos para actualizar state no fim
all_items.extend(items)
page += 1
# prevenção simples: se a API não suportar paginação por updated_after, podes precisar de parar por segurança
if page > 1000:
break
time.sleep(0.2) # suavizar chamadas
if all_items:
new_state = process_and_update_state(conn, all_items)
print('State actualizado para', new_state)
else:
print('Sem novidades')
if __name__ == '__main__':
conn = init_state_db()
run_once(conn)
Verificar o resultado
Confirma que o campo state foi gravado na tabela etl_state (SELECT * FROM etl_state). Verifica também que só foram importados registos com updated_at > last_updated_at anterior. Faz duas execuções: a primeira deverá importar dados e gravar state; a segunda não deverá importar nada se não houver alterações.
Conclusão
Com um state simples e chamadas condicionadas podes transformar uma API de Dados num fluxo incremental eficiente: menos tráfego, ETL mais rápido e menor latência. Próximos passos: acrescentar retries com backoff exponencial, registo (logging) detalhado e testes automáticos. Dica: começa por validar o formato de timestamp da API para evitar erros de fusos horários e parsing.