Nuestro sistema de ingesta de logs llevaba 18 meses funcionando con polling cada 30 segundos. Funcionaba, en el sentido en que el avión vuela aunque el motor derecho eche humo. Cuando un pipeline de Airflow fallaba a las 2 de la madrugada, nuestros ingenieros lo descubrían a la 2:00:30 — o a la 2:01:00, dependiendo de dónde cayera el ciclo. Para debugging de producción eso es una eternidad.
Hace cuatro meses rediseñamos el sistema completo. Esta es la arquitectura que construimos, por qué descartamos la alternativa más obvia, y dónde todavía falla.
Por qué el polling cada 30s destruye el debugging de pipelines
El problema con el polling no es la latencia en sí. Es lo que esa latencia hace al proceso de debugging.
Cuando un pipeline de datos falla en el paso 3 de 8 y tardas 30 segundos en enterarte, los pasos 4-8 ya intentaron ejecutarse. Algunos fallaron silenciosamente. Algunos escribieron datos parciales en destino. Para entonces tienes cinco errores encadenados donde solo había uno real.
El p99 de detección era 28 segundos. En práctica, eso significaba que el primer error que veía el ingeniero de guardia era el quinto de la cadena, no el primero.
La arquitectura: Kafka + CDC con Debezium
La nueva pila tiene tres componentes principales: Debezium captura los cambios en Postgres a nivel de write-ahead log, los publica en un topic de Kafka, y nuestro consumer los procesa e indexa en tiempo real.
Debezium lee directamente el WAL de Postgres. No hace queries. No añade carga al servidor. Cuando un pipeline escribe una fila de log, Debezium la captura en el WAL y la publica en Kafka en milisegundos — independientemente de lo que esté haciendo el resto del sistema.
El consumer de Kafka en Starlogs procesa los eventos, los indexa, evalúa las reglas de alerta configuradas y dispara notificaciones. La latencia end-to-end — desde que el pipeline escribe el log hasta que el ingeniero recibe la alerta — es ahora 340ms en p99.
Por qué descartamos webhooks como alternativa
La solución más obvia era pedir a los clientes que configuraran webhooks desde sus pipelines hacia nuestra API. Sencillo de explicar, cero infraestructura nuestra.
Lo descartamos por tres razones concretas:
| Opción | Ventaja | Por qué la descartamos |
|---|---|---|
| Webhooks desde pipeline | Sin infraestructura propia | Requiere modificar el código del pipeline. Fallan silenciosamente si la red cae. No capturan errores que impiden que el código llegue a ejecutarse. |
| Polling cada 5 segundos | Cero cambios en arquitectura | Mejora la latencia solo 6×, no 82×. Multiplica la carga en Postgres por 6. No resuelve el problema de los errores encadenados. |
| Push desde el ORM | Granularidad perfecta | Acoplamiento fuerte al stack del cliente. Solo funciona si usan SQLAlchemy o Tortoise. La mitad de nuestros clientes no lo hacen. |
| CDC con Debezium + Kafka | Transparente para el cliente, latencia sub-segundo, resiliente a fallos de red | Complejidad operacional alta. Requiere acceso al WAL de Postgres (no disponible en todos los managed services). |
La decisión más difícil fue aceptar la complejidad operacional. Operar Debezium y Kafka añade superficie a mantener. Llevamos cuatro meses en producción y hemos tenido dos incidentes menores — ambos relacionados con replication slots de Postgres que se quedaron colgados. Nada que afectara a clientes, pero trabajo real.
Cómo inicializar el consumer en un pipeline de Airflow
La integración no requiere modificar el código del DAG. El consumer se inicializa a nivel de conexión de Postgres y captura todos los eventos de escritura automáticamente.
from starlogs import StarlogsHook from airflow.providers.postgres.hooks.postgres import PostgresHook # Inicializar una vez por DAG — NO por tarea # StarlogsHook envuelve la conexión y registra el pipeline_id # en el consumer de Kafka. Los logs fluyen automáticamente. starlogs = StarlogsHook( conn_id="postgres_warehouse", pipeline_id="etl_ventas_diario", env="production", # alert_on: qué niveles disparan notificación inmediata alert_on=["ERROR", "CRITICAL"], ) def extract_ventas(**context): pg = PostgresHook(postgres_conn_id="postgres_warehouse") # La conexión ya está instrumentada. Cada query queda registrada. rows = pg.get_records("SELECT * FROM ventas WHERE fecha = %s", parameters=[context["ds"]]) return rows def transform_ventas(rows, **context): # Si esta función lanza una excepción, Starlogs la captura # desde el WAL antes de que Airflow marque la tarea como failed. # Eso te da ~200ms de ventaja para correlacionar el error. result = [] for row in rows: result.append({ "producto": row[0], "importe_eur": float(row[1]), }) return result
StarlogsHook registra un replication slot en Debezium. Si inicializas uno por tarea, acabas con slots huérfanos que bloquean el WAL y pueden detener Postgres. Lo aprendimos a las 3 de la mañana de un martes.
Dónde todavía falla: el muro de los pipelines batch largos
CDC con Debezium funciona sobre el WAL de Postgres. El WAL tiene una retención configurable — por defecto, en la mayoría de managed services, es de 24 horas.
Para pipelines batch que duran más de 6 horas (tenemos clientes con ETLs de fin de semana que procesan 72 horas de datos), existe un riesgo real: si el consumer de Kafka se cae durante una ventana larga, el WAL puede rotar y perdemos eventos.
Ser directos sobre esto no es heroísmo editorial. Es pragmatismo. Un cliente que descubre esta limitación en producción a las 3am nos va a churnar.
Lo que no funcionó la primera vez
Pasamos dos semanas con un bug extraño: los eventos llegaban duplicados al consumer, con un desfase de exactamente 30 segundos. El número nos sonó familiar.
Resulta que el sistema de polling anterior no había sido desactivado del todo. Seguía corriendo como cron job legacy en una instancia de ECS que habíamos "apagado" pero que tenía Auto Scaling habilitado. Cada vez que el sistema CDC generaba carga, el cron job del polling se despertaba y publicaba sus propios eventos al mismo topic de Kafka.
La solución fue obvia una vez que la vimos. Llegar hasta ahí llevó nueve sesiones de debugging y un grafo de correlación de eventos que todavía está en nuestra Notion como recordatorio de que "apagar" un servicio y "destruir" un servicio son cosas distintas.
¿Vale la complejidad?
Depende del caso de uso. Si tus pipelines corren en menos de 4 horas y operas tu propio Postgres con configuración de WAL controlada, sí, rotundamente. La reducción de 82× en latencia de detección se traduce en errores más simples, debugging más rápido, y noches más tranquilas.
Si tienes pipelines de fin de semana o usas un managed service con WAL de 24h y sin acceso a configuración — espera al modo híbrido de Q3, o usa el SDK directo con buffer.
El código del consumer y el conector de Debezium están en github.com/starlogs/kafka-cdc-consumer con MIT license. Hay un script de instalación que levanta el stack completo en Docker en 4 minutos, incluyendo una instancia de Postgres de prueba con datos sintéticos de pipelines.