Planteamiento y contexto
Los trabajos ascendentes producen varios activos de datos cada día. Los reportes y las comprobaciones de calidad deben ejecutarse después de que sus dependencias se actualicen. Explica cómo modelarías productores, consumidores, particiones y recuperación con la programación orientada a activos de Airflow, y en qué se diferencia de las programaciones temporales y los sensores externos.
Qué está evaluando el entrevistador
- Tratar un activo como una dependencia lógica de datos que actualiza una tarea, no solo como un nombre de archivo o una etiqueta cron.
- Distinguir eventos de activos, calendarios de DAG (timetables), composición de expresiones y disparadores de eventos.
- Considerar eventos duplicados, datos tardíos, granularidad de particiones, controles de calidad (quality gates) y backfills.
- Explicar monitoreo, autorización, idempotencia, reintentos y comportamiento de pausa/reanudación.
Preguntas aclaratorias para hacer
- ¿Un activo representa una tabla completa, una partición o una ruta de almacenamiento de objetos? ¿Los eventos llevan la identidad de partición y de lote?
- ¿Debe el consumidor esperar a cada activo ascendente o cualquier actualización individual puede dispararlo? ¿Se necesitan expresiones de activos entre distintos DAGs?
- Si llega un evento pero fallan las comprobaciones de calidad, ¿debe bloquearse al consumidor, reintentarse el productor o requerirse una liberación manual por un humano?
- ¿Necesitamos backfills históricos, reproducción de eventos o compatibilidad con una programación cron existente? ¿Cuál es el costo de un disparo duplicado?
Una respuesta de 30 segundos
Modelaría el producto de datos como activos con contratos de actualización, permitiendo luego que un DAG productor emita una actualización de activo solo después de que una escritura atómica y las comprobaciones de calidad tengan éxito. Un DAG consumidor utiliza dependencias o expresiones de activos para definir su condición de disparo en lugar de adivinar la disponibilidad con un intervalo cron más corto. Definiría reglas de partición, idempotencia, eventos duplicados, actualizaciones tardías y backfills, monitoreando luego la latencia de eventos, los DAGs en espera, los reintentos y la frescura. Si las actualizaciones se originan fuera de Airflow, evaluaría la recuperabilidad y los límites de seguridad de un disparador basado en eventos.
Análisis detallado paso a paso
1. Establecer un contrato de activo
La documentación oficial de Airflow Assets define los activos como dependencias de datos compartidas entre DAGs. Un productor debe actualizar un activo solo después de que el commit de datos y las comprobaciones de contrato tengan éxito; crear un archivo temporal, iniciar una tarea o escribir datos parcialmente no constituye disponibilidad. Mantén estables el URI del activo, el propietario y la granularidad de partición para que los consumidores puedan distinguir los datos nuevos de los antiguos.
2. Seleccionar la lógica de disparo
Un DAG consumidor puede depender de uno o más activos. La documentación oficial de Asset-Aware Scheduling admite combinaciones lógicas que expresan condiciones como todos los activos actualizados o cualquier activo actualizado; un timetable puede permanecer como una restricción independiente. Define qué sucede cuando las condiciones de evento y de tiempo coinciden, y no confundas una expresión con un filtro sobre el contenido de las filas.
3. Manejar particiones y duplicados
Un evento de activo no reemplaza una marca de agua (watermark) de partición de datos. Lleva una identidad de lote o partición trazable y utiliza una tabla de marcas de agua, una clave única o una escritura transaccional para garantizar la idempotencia. Los eventos duplicados, los reintentos de tareas y la recuperación del programador pueden volver a evaluar las dependencias, por lo que los consumidores deben ser ejecutables nuevamente de forma segura en lugar de asumir exactamente un solo evento.
4. Eventos externos, calidad y recuperación
Cuando las actualizaciones provienen de una cola u otro sistema, utiliza la documentación oficial de programación basada en eventos para elegir un disparador, verificando luego las comprobaciones recuperables, las credenciales de conexión y la liberación de recursos. Mantén un activo como no listo cuando la calidad falle, y reintenta o requiere una liberación explícita. Para los backfills, decide si emitir eventos, aislar particiones históricas y evitar que una reproducción sobrescriba el resultado en tiempo real.
Respuesta modelo
Primero definiría el contrato del activo: URI, clave de partición, productor, control de calidad e identidad de lote. Un DAG productor actualiza el activo solo después de una escritura atómica y comprobaciones de calidad. Un DAG consumidor utiliza dependencias o expresiones de activos para condiciones de todos-los-ascendentes frente a cualquier-ascendente, opcionalmente combinadas con un timetable. Los eventos llevan la identidad de partición y de lote; el consumidor utiliza marcas de agua y claves de idempotencia para duplicados, reintentos y recuperación del programador. Las actualizaciones externas utilizan un disparador controlado con credenciales acotadas y tiempo de vida de recursos limitado. Monitoreo la latencia de eventos, las tareas en espera, la frescura y los fallos de reproducción. Las rutas de backfill y de tiempo real utilizan límites de lote separados para que la reproducción histórica no pueda sobrescribir el resultado más reciente.
Errores comunes
- Tratar un activo como una ruta arbitraria sin definir la propiedad de la actualización ni la disponibilidad.
- Acortar un intervalo cron sin una señal de datos listos, un contrato de partición o un control de calidad.
- Asumir que los eventos de activos se entregan exactamente una vez e ignorar reintentos, duplicados o la recuperación del programador.
- Tratar una expresión de activos como un filtro de contenido de filas y omitir la marca de agua de partición real.
- Permitir que un disparador externo retenga conexiones o credenciales indefinidamente sin tiempo de espera (timeout) ni limpieza.
- Reutilizar el DAG de tiempo real para backfills y permitir que la reproducción histórica sobrescriba la salida en vivo.
Preguntas de seguimiento y respuestas
¿Qué pasa si dos activos ascendentes llegan en momentos diferentes?
Decide si el consumidor espera a todos los activos o si puede publicar un resultado parcial. Para la semántica de todos los activos, rastrea la marca de agua de llegada de cada partición y genera alertas ante tiempos de espera agotados. Para resultados parciales, publica una versión de modo que los consumidores sepan que un activo posterior podría actualizarlo.
¿Cómo probarías eventos duplicados?
En un entorno de pruebas, publica la misma actualización de activo dos veces, reintenta el productor y reinicia el programador. Verifica las claves únicas del consumidor, las marcas de agua y los límites de transacción, incluidos los recuentos de filas, las versiones y los efectos secundarios externos como notificaciones o cobros.
¿Cuándo mantendrías un sensor externo?
Mantén uno temporalmente cuando el sistema dependiente no pueda emitir eventos de activos de Airflow, solo exponga una interfaz de sondeo (polling) controlada o deba mantener compatibilidad durante una migración. Registra la fecha límite de migración y el costo del sondeo, avanzando luego el sistema ascendente hacia una señal de actualización de activo verificable.