Tienes el almacén del capítulo 15 y la carga incremental del capítulo 16. Falta la parte que nadie enseña y que es la que suena el teléfono: qué pasa cuando eso corre solo a las tres de la mañana y algo sale mal 🌙
Tres preguntas, y las tres tienen que tener respuesta antes de programar nada:
- ¿Alguien se entera de que falló, o se descubre el lunes en una reunión?
- ¿Qué quedó a medias? ¿El tablero está viejo, vacío o mentiroso?
- ¿Cómo se vuelve a correr sin romper lo que sí funcionó?
Un almacén para trastear
import os import shutil import sqlite3 import tempfile carpeta = tempfile.mkdtemp() almacen = os.path.join(carpeta, 'almacen.db') shutil.copy('tienda.db', almacen) con = sqlite3.connect(almacen) con.execute("""CREATE TABLE carga_log ( paso TEXT, corrio_en TEXT, filas INTEGER, estado TEXT, detalle TEXT)""") con.commit() print('almacen listo, con', con.execute('SELECT count(*) FROM pedidos').fetchone()[0], 'pedidos')
almacen listo, con 900 pedidos
El orquestador entero, en dieciocho líneas
Lo escribo a mano porque cuando veas Airflow quiero que reconozcas lo que hace por dentro. Un paso es un nombre, de qué otros pasos depende, y una función:
def anota(paso, filas, estado, detalle=''): con.execute('INSERT INTO carga_log VALUES (?, ?, ?, ?, ?)', (paso, '2026-08-21 03:00', filas, estado, detalle)) con.commit() def corre(pasos): hechos = set() for nombre, necesita, funcion in pasos: if not necesita.issubset(hechos): faltan = ', '.join(sorted(necesita - hechos)) anota(nombre, 0, 'saltado', 'falta ' + faltan) print(f'{nombre:9} saltado, falta {faltan}') continue try: filas = funcion() except Exception as e: anota(nombre, 0, 'error', str(e)) print(f'{nombre:9} ERROR {e}') continue hechos.add(nombre) anota(nombre, filas, 'ok') print(f'{nombre:9} ok {filas} filas') print('el motor son 18 lineas')
el motor son 18 lineas
Fíjate en lo que hace y en lo que no hace. Si un paso falla, los que dependen de él no corren y quedan anotados como saltados. No se para todo de golpe ni se sigue como si nada: se sigue con lo que no dependía del que falló.
Eso es una de las palabras que vas a oír, el DAG. Es un dibujo de qué necesita qué, y la letra que importa es la A de acíclico: un paso no puede depender de sí mismo dando la vuelta, porque entonces no hay por dónde empezar.
Los cuatro pasos
def extrae(): con.execute('DROP TABLE IF EXISTS cruda_pedidos') con.execute("""CREATE TABLE cruda_pedidos AS SELECT *, '2026-08-21' AS cargado_en FROM pedidos""") return con.execute('SELECT count(*) FROM cruda_pedidos').fetchone()[0] def limpia(): con.execute('DROP TABLE IF EXISTS limpia_pedidos') con.execute("""CREATE TABLE limpia_pedidos AS SELECT id, fecha, trim(canal) AS canal, monto, cargado_en FROM cruda_pedidos WHERE monto IS NOT NULL AND monto > 0""") return con.execute('SELECT count(*) FROM limpia_pedidos').fetchone()[0] def revisa(): entraron, quedaron = con.execute( 'SELECT (SELECT count(*) FROM cruda_pedidos), ' ' (SELECT count(*) FROM limpia_pedidos)').fetchone() perdidas = entraron - quedaron if perdidas > entraron * 0.05: raise ValueError(f'se cayeron {perdidas} de {entraron} pedidos al limpiar') return quedaron def publica(): con.execute('DROP TABLE IF EXISTS ventas_dw') con.execute("""CREATE TABLE ventas_dw AS SELECT canal, substr(fecha, 1, 7) AS mes, count(*) AS pedidos, round(sum(monto), 2) AS venta FROM limpia_pedidos GROUP BY canal, mes""") return con.execute('SELECT count(*) FROM ventas_dw').fetchone()[0] PASOS = [('extrae', set(), extrae), ('limpia', {'extrae'}, limpia), ('revisa', {'limpia'}, revisa), ('publica', {'revisa'}, publica)] corre(PASOS)
extrae ok 900 filas limpia ok 873 filas revisa ok 873 filas publica ok 72 filas
Cuatro en verde. Y mira revisa, que es el paso que casi nadie
escribe: no transforma nada. Solo mira si el resultado del paso
anterior tiene sentido, y si no lo tiene, revienta a propósito.
Se cayeron 27 pedidos al limpiar, de 900. Son los de monto nulo o cero, que ya estaban ahí. Menos del 5% que puse de tope, así que la carga sigue.
print(con.execute('SELECT round(sum(venta), 2) FROM ventas_dw').fetchone()[0], 'soles en el tablero')
532653.85 soles en el tablero
Ahora que salga mal
Se rompe algo en el sistema de origen y una de cada siete ventas empieza a llegar con monto negativo. No es raro: pasa cuando alguien cambia cómo se guardan las devoluciones y no avisa.
con.execute('UPDATE pedidos SET monto = -1 WHERE id % 7 = 0')
con.commit()
corre(PASOS)
extrae ok 900 filas limpia ok 747 filas revisa ERROR se cayeron 153 de 900 pedidos al limpiar publica saltado, falta revisa
Ahí está el capítulo entero en cuatro líneas.
El paso de limpieza no falló: hizo su trabajo y devolvió 747 filas tan contento, porque descartar filas que no cumplen es literalmente lo que le pediste. Quien avisa es el chequeo, comparando contra lo que entró.
Y lo importante es lo que no pasó:
print(con.execute('SELECT round(sum(venta), 2) FROM ventas_dw').fetchone()[0], 'soles en el tablero')
532653.85 soles en el tablero
El tablero sigue con los números de la carga buena. Está viejo, y está bien.
Si publica hubiera corrido, la jefa de ventas abriría el lunes un
tablero fresquito al que le faltan 153 pedidos, sin ninguna señal de que le
falta nada 😐
Viejo y correcto se explica en una frase. Nuevo y a medias no se explica nunca, porque para explicarlo primero hay que darse cuenta, y nadie se da cuenta 💛
Un chequeo, visto de cerca
Llama al chequeo tú misma, sin el orquestador que lo atrapa:
revisa()
ValueError: se cayeron 153 de 900 pedidos al limpiar
Eso es todo lo que es un chequeo de calidad: una consulta y un
raise. Las herramientas caras traen cientos escritos, y el que te
va a salvar es el que escribas tú, porque conoce tu negocio.
Los cuatro que pondría siempre:
- Llegaron filas. Cero filas casi nunca es un día tranquilo.
- No se perdieron por el camino. El cuadre del capítulo 15, como paso.
- La clave no se repite. Un
count(*)contra uncount(DISTINCT id). - Los números están en su rango. Ninguna venta negativa, ninguna fecha del futuro.
Y la bitácora, para el lunes
for fila in con.execute('SELECT paso, filas, estado, detalle FROM carga_log ' 'ORDER BY rowid DESC LIMIT 4').fetchall(): print('%-9s %4d %-8s %s' % fila)
publica 0 saltado falta revisa revisa 0 error se cayeron 153 de 900 pedidos al limpiar limpia 747 ok extrae 900 ok
Cuatro líneas que contestan qué pasó, en qué paso y por qué, sin que nadie tenga que acordarse. Es la misma bitácora del capítulo anterior con una columna más.
Las herramientas, para que no te vendan humo
| Herramienta | Qué agrega sobre lo de arriba | Cuándo |
|---|---|---|
| cron | lo dispara a una hora, y nada más | una o dos cargas simples |
| Airflow | el DAG con pantalla, reintentos, backfill, historial | muchas cargas que dependen entre sí |
| Dagster, Prefect | lo mismo con otra filosofía y menos ceremonia | igual, si empiezas de cero hoy |
| dbt | ordena y prueba las transformaciones en SQL, no las mueve | cuando la parte difícil ya es el SQL |
Ninguna de las cuatro te va a decir qué chequeo poner ni cuál es el orden correcto. Eso lo pones tú, y por eso vale la pena entenderlo con dieciocho líneas antes de instalar nada 🧾
Programar que una consulta corra sola cada noche
| PostgreSQL | no trae planificador: se usa la extension pg_cron o el cron del sistema |
| MySQL | CREATE EVENT ... ON SCHEDULE EVERY 1 DAY, si event_scheduler esta en ON |
| SQL Server | SQL Server Agent, que es un servicio aparte y hay que tenerlo instalado |
| SQLite | nada: no hay ningun proceso corriendo al que programarle algo |
Y de los cuatro, el que mas se usa en proyectos de verdad es ninguno: el planificador vive fuera de la base, porque una carga casi nunca es solo SQL. Baja un archivo, llama a una API, escribe en otro sitio. La base es un paso, no el director.
Comprueba que lo tienes
La carga de anoche se cayó en el paso de limpieza y el paso que publica el tablero no llegó a correr. ¿Qué es lo mejor que puede pasar a la mañana siguiente?
- Que el tablero enseñe los números de anteayer y la bitácora diga que la carga falló
- Que el tablero salga vacío para que se note
- Que el tablero publique lo que alcanzó a limpiarse
- Que la carga se reintente sola hasta que pase
Ejercicios
1. Reintentar lo que se puede reintentar
El servidor de origen se cae a ratos. Escribe el reintento con espera creciente.
import time intentos = {'n': 0} def baja_del_sistema(): intentos['n'] += 1 if intentos['n'] < 3: raise ConnectionError('timeout contra el servidor de pedidos') return 900 def con_reintentos(funcion, veces=4, espera=0.01): for intento in range(1, veces + 1): try: return funcion() except ConnectionError as e: print(f'intento {intento}: {e}') if intento == veces: raise time.sleep(espera * 2 ** (intento - 1)) print('filas:', con_reintentos(baja_del_sistema))
intento 1: timeout contra el servidor de pedidos intento 2: timeout contra el servidor de pedidos filas: 900
La espera que se dobla en cada vuelta tiene nombre, espera exponencial, y existe porque el origen suele estar caído justo porque está saturado. Reintentar cada segundo es echarle más leña.
Y mira qué error atrapa: ConnectionError y nada más. Un reintento
que atrapa todo se traga también el error de datos, y entonces reintenta cuatro
veces algo que iba a fallar igual las cuatro 🧾
2. Volver a correr días viejos
Escribe una carga por día que se pueda relanzar para cualquier fecha. Arriba le metimos montos negativos a la base a propósito, así que empiezo copiándola otra vez limpia.
limpio = os.path.join(carpeta, 'limpio.db') shutil.copy('tienda.db', limpio) con = sqlite3.connect(limpio) def carga_del_dia(dia): con.execute('DELETE FROM ventas_dia WHERE fecha = ?', (dia,)) con.execute("""INSERT INTO ventas_dia SELECT fecha, count(*), round(sum(monto), 2) FROM pedidos WHERE fecha = ? GROUP BY fecha""", (dia,)) con.commit() return con.execute('SELECT count(*) FROM ventas_dia WHERE fecha = ?', (dia,)).fetchone()[0] con.execute('CREATE TABLE ventas_dia (fecha TEXT, pedidos INTEGER, venta REAL)') for dia in ('2026-06-20', '2026-06-21', '2026-06-22', '2026-06-23'): carga_del_dia(dia) for fila in con.execute('SELECT * FROM ventas_dia ORDER BY fecha').fetchall(): print(fila)
('2026-06-20', 2, 1224.52)
('2026-06-21', 2, 1450.95)
('2026-06-22', 2, 1806.89)
('2026-06-23', 2, 1541.36)
Que la carga reciba el día como parámetro es lo que hace posible el backfill: recorrer con un bucle los días que faltan o que salieron mal.
Una carga escrita como "trae lo de ayer" no se puede rehacer para el martes pasado sin cambiarle el código, y ese es el momento en que alguien acaba corrigiendo el almacén a mano 🙃
3. Comprueba que el backfill no ensucia
Relanza dos días que ya estaban cargados.
for dia in ('2026-06-20', '2026-06-21'): carga_del_dia(dia) print(con.execute('SELECT count(*), round(sum(venta), 2) FROM ventas_dia').fetchone())
(4, 6023.72)
Cuatro filas y la misma venta que antes. El DELETE de la primera
línea es lo que lo consigue, y es la versión más simple de idempotencia que
existe: borra lo tuyo y vuelve a ponerlo.
Fíjate en que borra solo ese día. Un backfill que empieza por
DELETE FROM ventas_dia a secas arregla el martes y te borra los
otros tres años.
4. El paso que no se puede reintentar
Un acumulador. Córrelo dos veces y mira qué pasa.
con.execute('CREATE TABLE saldo_canal (canal TEXT PRIMARY KEY, acumulado REAL)') con.execute("INSERT INTO saldo_canal SELECT canal, 0 FROM pedidos GROUP BY canal") def suma_el_dia(dia): con.execute("""UPDATE saldo_canal SET acumulado = acumulado + coalesce( (SELECT sum(monto) FROM pedidos p WHERE p.canal = saldo_canal.canal AND p.fecha = ?), 0)""", (dia,)) con.commit() suma_el_dia('2026-06-20') print('una vez :', con.execute('SELECT round(sum(acumulado), 2) FROM saldo_canal').fetchone()[0]) suma_el_dia('2026-06-20') print('dos veces:', con.execute('SELECT round(sum(acumulado), 2) FROM saldo_canal').fetchone()[0])
una vez : 1224.52 dos veces: 2449.04
El día se sumó dos veces y el saldo salió al doble. Este paso no se puede reintentar, y ponerle reintentos automáticos es peor que no ponérselos.
Se arregla de dos maneras. La buena: no acumular, recalcular desde cero cada vez, que es lo que hacía el ejercicio 2. La otra: anotar en la bitácora qué días ya se sumaron y comprobarlo antes.
La regla de dedo que uso: si un paso empieza con la palabra "sumar" o "añadir", desconfía. Si empieza con "recalcular", duerme tranquila 💛
5. Cuándo suena el teléfono
Tienes la bitácora. ¿Qué merece despertar a alguien a las tres de la mañana y qué espera al desayuno?
Un paso en error que deja el tablero viejo: espera. El daño ya está contenido, el tablero de ayer sigue siendo correcto, y a las siete se arregla igual de bien que a las tres.
Una carga que terminó en verde con cero filas tres días seguidos: esto sí. No falló nada y por eso nadie miró, que es exactamente cómo un origen deja de mandar datos durante una semana sin que nadie lo note.
Un chequeo de duplicados que salta: depende de si publicó. Si el chequeo va antes de publicar, espera. Si va después, es urgente, y eso ya te dice que estaba en el sitio equivocado.
La pregunta que ordena todo esto es una sola: ¿alguien va a tomar una decisión con este dato antes de que yo llegue? Si la respuesta es no, no hay urgencia, y despertar a la gente por cosas que podían esperar es la forma más rápida de que dejen de mirar las alertas 🚨
6. Ordena estos seis pasos
Una carga real, desordenada. Di qué depende de qué.
Los seis: publicar el tablero, bajar el CSV del marketplace, comprobar que no hay pedidos duplicados, cargar la tabla cruda, actualizar el maestro de clientes, armar la tabla de hechos.
Bajar el CSV va primero y no depende de nada. La tabla cruda depende del CSV. Actualizar el maestro de clientes también depende del CSV, y acá está lo interesante: esos dos pueden correr a la vez, porque ninguno necesita al otro.
La tabla de hechos necesita las dos, porque cruza ventas con clientes. El chequeo de duplicados va después de la tabla de hechos. Y publicar va al final, después del chequeo.
Eso último es lo único que no se negocia. Si el chequeo va después de publicar, no es un chequeo: es una autopsia 🙃
Preguntas frecuentes
¿Qué es un DAG?
El dibujo de qué paso necesita a qué paso. La A es de acíclico: un paso no puede depender de sí mismo dando la vuelta, porque entonces no habría por dónde empezar.
¿Qué es dbt?
Una herramienta para ordenar y probar transformaciones escritas en SQL. No mueve datos: los transforma donde ya están, y su aporte real son las pruebas y la documentación.
¿Qué es Airflow?
Un orquestador: dispara los pasos en el orden del DAG, reintenta lo que falla y guarda el historial. Por dentro hace lo mismo que las dieciocho líneas de este capítulo, con pantalla.
¿Qué es un backfill?
Volver a correr la carga para fechas pasadas. Solo se puede si la carga recibe el día como parámetro, y por eso se escribe así desde el principio.
Practica este capítulo 📓
Todo el código de arriba en un cuaderno que corre de principio a fin, y los ejercicios con una celda vacía para que los hagas tú. Se abre en Google Colab de un clic y no hay que instalar nada. Donde veas %%revisa, escribe tu respuesta y el cuaderno te dice si te salió.
¿Prefieres trabajar en tu máquina? Bájate el cuaderno de práctica o el de soluciones. Todos están también en github.com/soymissyera/MissYeraEjercicios.