Cómo extraer cambios incrementales en APIs de Datos: paso a paso
Aprende a extraer solo los registros nuevos o alterados de una API de Datos — útil para reducir tráfico, acelerar ETL y mantener una copia local actualizada. Vamos a explicar el porqué y mostrar un ejemplo en Python que detecta y procesa cambios incrementales.
Prerequisitos
- Python 3.8+ instalado
- Bibliotecas: requests y sqlite3 (integrada en Python)
- Una API REST que exponga un campo de fecha/timestamp o un campo incremental (p. ej.: updated_at o id)
- Nociones básicas de HTTP y JSON
Paso 1: Entender el modelo de cambio de la API
Antes de escribir código, verifica si la API tiene un campo que indique cambios (por ejemplo updated_at) o un identificador incremental. Si tiene updated_at en UTC, puedes pedir solo registros posteriores a una fecha. Si solo hay id incremental, puedes pedir id > último_id. Saber esto evita tirar de todos los datos y es la base de la extracción incremental.
Paso 2: Estructurar un almacenamiento de control (state)
Necesitas guardar el último timestamp o último id extraído. Un archivo simple o una pequeña base de datos SQLite funciona bien para un flujo ETL simple. Aquí usamos SQLite para persistencia 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()
Paso 3: Hacer una llamada condicionada a la API
Construye la query usando el valor guardado. Ejemplo con parámetro updated_after (muchas APIs aceptan algo similar). Maneja errores comunes como 400/500 y limitaciones de rate limit; si falta soporte para filtrado, tendrás que usar paginación y 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 # nombre del parámetro depende de la API
resp = requests.get(API_URL, params=params, timeout=10)
resp.raise_for_status()
return resp.json() # asume JSON con lista y metadatos
Paso 4: Procesar resultados y actualizar el state
Al recibir los registros, procésalos (grabar en tu base local, transformar, etc.) y calcula el nuevo valor de state: max(updated_at) o max(id). Solo después de guardar con éxito actualiza el state para evitar pérdida de datos.
def process_and_update_state(conn, items):
# Ejemplo mínimo: grabar en archivo o base local (omito implementación completa de escritura)
# Supone que cada item tiene 'id' y 'updated_at' en ISO 8601
if not items:
return None
max_ts = max(item['updated_at'] for item in items)
# aquí escribiríamos los items en una tabla o archivos
set_state(conn, 'last_updated_at', max_ts)
return max_ts
Paso 5: Juntar todo en un ciclo seguro
Ejecuta un ciclo que lea el state, llame a la API por páginas hasta que no haya más, procese y actualice el state al final. Maneja errores temporales con reintentos simples y respeta rate limits con backoff lineal o 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('Error HTTP:', e)
break
items = data.get('items', [])
if not items:
break
# lógica de procesamiento aquí; acumulamos para actualizar state al final
all_items.extend(items)
page += 1
# prevención simple: si la API no soporta paginación por updated_after, puede que necesites parar por seguridad
if page > 1000:
break
time.sleep(0.2) # suavizar llamadas
if all_items:
new_state = process_and_update_state(conn, all_items)
print('State actualizado a', new_state)
else:
print('Sin novedades')
if __name__ == '__main__':
conn = init_state_db()
run_once(conn)
Verificar el resultado
Confirma que el campo state se ha guardado en la tabla etl_state (SELECT * FROM etl_state). Verifica también que solo se importaron registros con updated_at > last_updated_at anterior. Haz dos ejecuciones: la primera debería importar datos y guardar state; la segunda no debería importar nada si no hay cambios.
Conclusión
Con un state simple y llamadas condicionadas puedes convertir una API de Datos en un flujo incremental eficiente: menos tráfico, ETL más rápido y menor latencia. Próximos pasos: añadir retries con backoff exponencial, registro (logging) detallado y tests automatizados. Consejo: comienza validando el formato de timestamp de la API para evitar errores de zonas horarias y parsing.