Capítulo 17 de 23 9 secciones 13 min

Compartir

Las tres de la mañana y nadie mirando

Pasos que dependen unos de otros, chequeos que paran la carga antes de publicar, reintentos y bitácora. Con un orquestador de dieciocho líneas para entender qué hacen los grandes.

Orquestar es decidir en qué orden corren los pasos de una carga, qué pasa cuando uno falla y cómo se vuelve a correr. Lo esencial son cuatro cosas: dependencias, chequeos que paran la carga antes de publicar, reintentos solo para lo que se puede repetir, y una bitácora que conteste si corrió. Airflow y compañía son eso con pantalla 🧾

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 un count(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

HerramientaQué agrega sobre lo de arribaCuándo
cronlo dispara a una hora, y nada másuna o dos cargas simples
Airflowel DAG con pantalla, reintentos, backfill, historialmuchas cargas que dependen entre sí
Dagster, Prefectlo mismo con otra filosofía y menos ceremoniaigual, si empiezas de cero hoy
dbtordena y prueba las transformaciones en SQL, no las muevecuando 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

PostgreSQLno trae planificador: se usa la extension pg_cron o el cron del sistema
MySQLCREATE EVENT ... ON SCHEDULE EVERY 1 DAY, si event_scheduler esta en ON
SQL ServerSQL Server Agent, que es un servicio aparte y hay que tenerlo instalado
SQLitenada: 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.

¿Le sirve a alguien que conoces?

Pásale el libro. Es gratis, está entero y no pide registro 🐣

Instagram y TikTok no dejan compartir enlaces desde la web: esos dos copian la URL para que la pegues en tu historia.

¿Tienes alguna duda o consulta?