Logo GCP con Eduardo

GCP con Eduardo

Descarga el código de la lección

Orquestación en Google Cloud: Triggers y Schedules en Data Fusion
✦ Guía Técnica & Arquitectura

Orquestación en Google Cloud: Triggers y Schedules en Data Fusion

Automatizar un pipeline ETL no significa simplemente poner un cron delante de un proceso. En una arquitectura empresarial, la verdadera pregunta es cómo controlar cuándo se ejecuta un pipeline, qué dependencias deben cumplirse, cómo transportar contexto entre ejecuciones y cómo evitar costes y ejecuciones duplicadas. Esta guía presenta un enfoque de nivel Senior / Staff / Architect para diseñar orquestación de pipelines GCP con Cloud Data Fusion.

Lo que aprenderás
Diseñar schedules robustos con frecuencias y expresiones cron.
Encadenar pipelines mediante triggers y dependencias explícitas.
Evitar solapamientos, ejecuciones duplicadas y problemas de FinOps.
Elegir entre triggers nativos y una capa de orquestación centralizada.
Google Cloud Cloud Data Fusion Apache Spark Cron ETL Triggers Managed Airflow Cloud Logging

1. El modelo mental correcto para la orquestación

Cloud Data Fusion permite construir pipelines como DAGs de etapas de integración, transformación y carga. Sin embargo, el DAG interno de un pipeline y el workflow que coordina varios pipelines son dos niveles diferentes de diseño.

El primer nivel resuelve el flujo de datos dentro del pipeline. El segundo resuelve el workflow empresarial: qué pipeline debe ejecutarse primero, qué condición habilita al siguiente, qué parámetros deben propagarse y qué sucede cuando una ejecución falla.

01 / DATA FLOW

DAG interno

Modela cómo los datos pasan desde sources a transformaciones, agregaciones y sinks. Es el nivel de procesamiento.

02 / PIPELINE FLOW

Triggers

Modelan dependencias entre pipelines. Un pipeline downstream puede reaccionar a la finalización de uno o varios pipelines upstream.

03 / WORKFLOW

Orquestador

Una capa como Managed Airflow resulta apropiada cuando las dependencias, observabilidad y coordinación cruzan múltiples pipelines y servicios.

2. Matriz de decisión: Schedule, Trigger o Airflow

Una decisión arquitectónica frecuente consiste en intentar resolver toda la orquestación exclusivamente mediante cron. Esto funciona mientras el workflow sea temporal y sencillo. Cuando aparecen dependencias de estado, fan-in, fan-out, branching o servicios externos, el modelo cambia.

Patrón Disparador Latencia conceptual Consistencia / dependencia Coste operativo Uso recomendado
Schedule Tiempo Programada Temporal Bajo ETL batch diario, horario o periódico.
Trigger Finalización de pipeline Orientada a evento Explícita entre upstream y downstream Bajo Encadenamientos sencillos entre pipelines Data Fusion.
Managed Airflow Tiempo, evento o lógica del DAG Controlada por workflow Alta expresividad Mayor Orquestación empresarial, dependencias complejas y cross-service.
Aplicación externa API / evento / sistema externo Variable Depende de la implementación Variable Casos donde el workflow pertenece a una plataforma externa.

3. Cómo programar un pipeline Data Fusion con Cron

Un schedule responde a una pregunta sencilla: “¿cuándo debo lanzar este pipeline?”

Para utilizarlo, el pipeline debe estar desplegado. Desde la vista del pipeline se configura el schedule y se puede elegir entre una configuración básica o una expresión cron avanzada.

Ejemplo: pipeline diario

Supongamos un pipeline denominado customer_daily_load que debe ejecutarse todos los días a las 02:00 UTC.

# Minuto Hora Día Mes DíaSemana
0      2    *   *   *

En producción, el detalle importante no es solamente la expresión cron. Debes analizar cuánto tarda el pipeline y qué ocurre si una ejecución todavía está activa cuando llega la siguiente ventana.

4. Pipeline Triggers: dependencias reales entre pipelines

Un trigger permite transformar una relación temporal en una relación de dependencia. En lugar de asumir que Pipeline B estará preparado a las 03:00 porque Pipeline A empezó a las 02:00, declaramos explícitamente: B depende de la finalización de A.

UPSTREAM

01 · Ingesta

Extrae datos desde las fuentes y los deja disponibles en una zona de staging.

DOWNSTREAM

02 · Transformación

Se inicia cuando la ingesta alcanza la condición de finalización definida.

CONSUMPTION

03 · Serving

Publica los datos preparados para consumo analítico o aplicaciones downstream.

La ventaja arquitectónica es fundamental: la dependencia deja de estar codificada como una suposición de horario y pasa a formar parte del workflow.

5. Antipatrones y errores críticos de producción

⚠ Antipatrón · Cron como mecanismo de dependencia

“A las 03:00 seguro que el pipeline anterior ya terminó”

Este diseño parece funcionar durante meses hasta que aumenta el volumen de datos, cambia la duración de Spark o aparece una degradación temporal. Entonces el pipeline downstream comienza mientras el dataset upstream todavía está incompleto.

El resultado puede ser peor que un fallo explícito: un pipeline aparentemente exitoso que consume datos parciales.

Solución: utiliza un trigger cuando la relación sea realmente una dependencia entre pipelines. Reserva el schedule para condiciones temporales.

Antipatrón FinOps: concurrencia sin límites

Un schedule frecuente puede producir ejecuciones simultáneas si la duración real del procesamiento supera el intervalo configurado. Además de aumentar el consumo, puede generar presión sobre fuentes, sinks y recursos de Spark.

Diseña explícitamente la política de concurrencia. Para cargas batch críticas, mide duración p50/p95, establece una ventana operacional y evita configurar frecuencias más agresivas que la capacidad real del pipeline.

Antipatrón de observabilidad: “si terminó, está bien”

El estado de ejecución no sustituye a las métricas de calidad. Un pipeline puede terminar correctamente y producir cero registros, duplicados o una cantidad anómala de datos.

Instrumenta volumen, errores, throughput y duración. Para pipelines críticos, combina observabilidad técnica con controles de calidad de datos.

6. Implementación práctica de un workflow parametrizado

Un patrón robusto consiste en evitar pipelines rígidos y utilizar argumentos de runtime para transportar contexto entre ejecuciones. De esta forma, el mismo pipeline puede procesar diferentes ventanas temporales sin modificar su diseño.

bash / conceptual deployment workflow production pattern
# -------------------------------------------------------------------
# 1. Variables del workflow
# -------------------------------------------------------------------

PROJECT_ID="my-gcp-project"
REGION="europe-west1"
NAMESPACE="production"
UPSTREAM="customer_ingestion"
DOWNSTREAM="customer_transform"

RUN_DATE="2026-08-28"

# -------------------------------------------------------------------
# 2. El pipeline upstream recibe contexto como runtime argument
# -------------------------------------------------------------------

gcloud data-fusion pipelines run \
  --project="${PROJECT_ID}" \
  --location="${REGION}" \
  --namespace="${NAMESPACE}" \
  "${UPSTREAM}" \
  --runtime-args="run_date=${RUN_DATE}"

# -------------------------------------------------------------------
# 3. El trigger del downstream debe transportar el mismo contexto
# -------------------------------------------------------------------
#
# Conceptualmente:
#
# customer_ingestion
#       |
#       | SUCCESS
#       v
# customer_transform
#       |
#       | runtime argument
#       v
# run_date=2026-08-28
#
# El pipeline downstream utiliza el argumento para seleccionar
# la partición o ventana temporal correspondiente.

El objetivo del ejemplo no es convertir la API de ejecución en el orquestador, sino mostrar un principio de diseño: el workflow debe transportar explícitamente el contexto de ejecución.

7. Payload configuration: pasar contexto entre pipelines

Una dependencia útil no solo comunica “A terminó”; también puede transportar información generada durante la ejecución upstream.

01

Runtime arguments

Fecha de procesamiento, partición, tenant, versión del dataset o identificador de ejecución.

02

Plugin configuration

Valores producidos por determinados plugins upstream que resultan necesarios para configurar el downstream.

03

Ventanas dinámicas

Permite reutilizar el mismo pipeline para diferentes horas, días, semanas o particiones sin duplicar workflows.

Regla de arquitectura

No hagas que el downstream “adivine” qué datos debe procesar. El contexto de ejecución debe ser explícito, trazable y reproducible.

8. ¿Cuándo pasar de Triggers a Managed Airflow?

Los triggers nativos son excelentes para dependencias sencillas. Sin embargo, una plataforma empresarial suele acabar necesitando coordinación entre múltiples sistemas: Cloud Data Fusion, BigQuery, Cloud Storage, APIs, validaciones, sensores, branching, retries y alertas.

TRIGGER

Usa triggers cuando...

Tienes una cadena sencilla de pipelines y la condición de ejecución depende principalmente del estado de uno o varios upstream.

AIRFLOW

Usa Managed Airflow cuando...

Necesitas centralizar DAGs, dependencias, monitorización, alertas y workflows que cruzan servicios o múltiples pipelines.

ARCHITECTURE

Evita mezclar sin criterio

Define una frontera clara. Un trigger local no debería convertirse gradualmente en una red difícil de razonar de dependencias distribuidas.

9. Patrones de diseño y mejores prácticas

01 · IDEMPOTENCIA

Diseña para reintentos

Una ejecución puede repetirse. Las escrituras deben ser idempotentes o utilizar mecanismos de deduplicación y particionamiento adecuados.

02 · PARAMETRIZACIÓN

Runtime arguments

Separa la lógica del pipeline de la ventana temporal y otros parámetros de ejecución.

03 · CONCURRENCIA

Controla overlap

Define qué sucede cuando una nueva ejecución llega mientras la anterior todavía está procesando datos.

04 · OBSERVABILIDAD

Mide el pipeline

Controla duración, registros de entrada y salida, errores y throughput para detectar regresiones de rendimiento.

05 · GOVERNANCE

Separación de entornos

Mantén convenciones claras para desarrollo, integración y producción, junto con IAM y control de cambios.

06 · FINOPS

Coste como métrica

La frecuencia del schedule debe justificarse por SLA y necesidad de negocio, no por el simple deseo de reducir la latencia.

10. Checklist de implementación empresarial

Clasifica el tipo de dependencia

Determina si el workflow depende del tiempo, de otro pipeline, de un evento externo o de una combinación de condiciones.

Despliega primero los pipelines

No diseñes la orquestación alrededor de pipelines que todavía no tienen una versión desplegable y observable.

Configura el schedule solo donde corresponda

Utiliza cron para ventanas temporales. Evita usar un horario artificial como sustituto de una dependencia real.

Modela las dependencias con triggers

Define upstream, downstream y las condiciones de finalización que deben activar la siguiente etapa.

Propaga el contexto de ejecución

Utiliza runtime arguments y payload configuration para que cada ejecución sepa exactamente qué ventana o dataset debe procesar.

Diseña la recuperación ante fallos

Define retries, comportamiento ante fallos, ejecución parcial, backfills y estrategia de replay antes de llegar a producción.

Controla concurrencia y coste

Mide duración y consumo. Asegura que la frecuencia de ejecución sea compatible con la capacidad disponible y los límites operativos.

Establece observabilidad

Centraliza logs y métricas, define alertas y monitoriza tanto fallos técnicos como anomalías de volumen de datos.

Revisa la complejidad del workflow

Si los triggers comienzan a formar una topología compleja entre numerosos sistemas, considera mover la coordinación a Managed Airflow.

11. Preguntas frecuentes sobre Data Fusion Triggers y Schedules

¿Cómo programar un pipeline de Cloud Data Fusion con cron?

Primero debes tener el pipeline desplegado. Desde la vista del pipeline puedes crear un schedule y utilizar la configuración avanzada para introducir una expresión cron. Para cargas batch, es importante considerar la duración del pipeline, la concurrencia y la zona horaria utilizada por la programación.

¿Cuál es la diferencia entre un schedule y un trigger?

Un schedule responde a una condición temporal: “ejecuta este pipeline a esta frecuencia”. Un trigger responde a una condición de workflow: “ejecuta este pipeline cuando otro pipeline termine según la condición definida”. En arquitectura de datos, esta diferencia evita convertir horarios aproximados en dependencias implícitas.

¿Cuándo conviene utilizar Managed Airflow?

Cuando el workflow requiere una coordinación más amplia que una dependencia sencilla entre pipelines: múltiples servicios, DAGs centralizados, dependencias complejas, branching, monitorización, alertas y control de ejecución. Para una cadena sencilla de pipelines Data Fusion, un trigger nativo suele ser una solución más simple.

Sobre el Autor: Eduardo Martínez Agrelo

AI & Data Architect

Eduardo Martínez Agrelo es AI & Data Architect, especializado en arquitectura de datos, inteligencia artificial y diseño de plataformas cloud. Su enfoque combina profundidad técnica, automatización, escalabilidad, gobernanza y optimización de costes para construir soluciones de datos preparadas para entornos empresariales.

En sus contenidos técnicos aborda decisiones arquitectónicas reales, patrones de producción y estrategias para llevar plataformas de datos e IA desde el prototipo hasta sistemas robustos y operables.