Google PDE: Orquestación de Flujos de Trabajo y Automatización de Canalizaciones — Guía de estudio
Forma parte de la Google Professional Data Engineer — Guía de estudio. Practica con respuestas verificadas en el centro de exámenes de Google, o realiza tests cronometrados en ExamRoll.io.
Descripción general
La orquestación de flujos de trabajo y la automatización de pipelines coordinan las tareas de datos entre servicios para que la ingesta, la transformación, los controles de calidad y la publicación se realicen de forma fiable, segura y rentable. En Google Cloud, la orquestación debe alinearse con el modelo de ejecución de cada carga de trabajo: lotes programados, flujos controlados por eventos, ad-hoc o trabajos de larga duración. Los objetivos de diseño son la repetibilidad, la idempotencia, la observabilidad, el privilegio mínimo y la promoción segura a través de los entornos.
Decisiones clave:
- Orquestación de lotes centrada en código con Cloud Composer (Apache Airflow) para DAGs, dependencias de tareas y programación avanzada.
- Coreografía de API sin servidor con Cloud Workflows para secuencias ligeras, controladas por eventos y entre servicios.
- Endpoints de ejecución como trabajos de Cloud Run o Dataproc, activados por Cloud Scheduler para cron o por Eventarc para eventos.
- Orquestación nativa de SQL con Dataform para transformaciones, aserciones y gestión de lanzamientos en BigQuery.
El modelo operativo pone énfasis en los reintentos con retroceso exponencial limitado, los tiempos de espera, los SLAs, las puestas al día (catchup) y los rellenos (backfills), el diseño de tareas idempotentes para reejecuciones seguras y un manejo robusto de fallos con captura en colas de mensajes fallidos (dead-letter). La seguridad se implementa mediante cuentas de servicio por pipeline, aislamiento de secretos, parametrización e IAM de privilegio mínimo. CI/CD, la infraestructura como código y una telemetría completa completan un enfoque listo para producción.
Orquestación en Google Cloud: Herramientas y Patrones
Cloud Composer (Airflow)
Los DAGs (Grafos Acíclicos Dirigidos) definen grafos de ejecución con dependencias explícitas. Usa la API TaskFlow u operadores (p. ej., BigQuery, Dataflow, Dataproc, Cloud Run) para expresar las tareas. Los sensores y los operadores diferibles reducen la carga del planificador (scheduler) para condiciones de espera (p. ej., la finalización de un objeto en Cloud Storage o la aparición de una partición en BigQuery).
Programación: las expresiones cron,
start_date,end_dateycatchupcontrolan las ejecuciones históricas. Usacatchuppara rellenos (backfills); deshabilítalo para destinos adyacentes a flujos (streaming) o no idempotentes. Limita la concurrencia conmax_active_runsy pools para proteger los sistemas dependientes (downstream).Dependencias:
set_upstream/set_downstreamo dependencias de taskflow. Para la orquestación basada en metadatos, genera tareas dinámicamente a partir de una tabla de control de BigQuery (p. ej., una lista de clientes/particiones) usando el mapeo dinámico de tareas, manteniendo estable el tiempo de análisis del DAG y haciendo que las tareas se basen en datos.Ejemplo de fragmento de DAG (conciso): from airflow import DAG from datetime import datetime, timedelta from airflow.providers.google.cloud.operators.dataflow import DataflowStartFlexTemplateOperator
default_args = dict(retries=3, retry_delay=timedelta(minutes=5), execution_timeout=timedelta(hours=2), sla=timedelta(hours=3)) with DAG(‘daily_csv_import’, start_date=datetime(2023,1,1), schedule_interval=‘0 2 * * *’, catchup=True, max_active_runs=1, default_args=default_args) as dag: import_job = DataflowStartFlexTemplateOperator( task_id=‘import’, body={’launchParameter’: {‘jobName’: ‘csv-import-{{ ds_nodash }}’, ‘parameters’: {‘dlq_table’: ‘bqproj.dlq.bad_rows’}}} )
Cloud Workflows, Cloud Scheduler, trabajos de Cloud Run y ejecución controlada por eventos
- Cloud Workflows orquesta las APIs de Google y los endpoints HTTP con reintentos incorporados, bucles, ramas paralelas y lógica de compensación. Es ideal para un flujo de control ligero entre servicios como BigQuery, Dataflow, Batch y trabajos de Cloud Run.
- Cloud Scheduler activa Workflows, temas de Pub/Sub o servicios HTTP para una automatización de estilo cron. Para un lote diario a las 02:00, programa un Workflow que lance un trabajo de Dataflow o Dataproc.
- Los trabajos de Cloud Run ejecutan pasos de lote en contenedores con reintentos automáticos y operaciones mínimas. Se combinan bien con Workflows para tareas de datos de varios pasos o para pre/post-procesamiento en torno a Dataflow o BigQuery.
- Controlado por eventos: usa Eventarc para enrutar eventos de finalización de objetos de Cloud Storage, mensajes de Pub/Sub o Audit Logs a Cloud Run o Workflows. Para notificaciones de trabajos de inserción de BigQuery en una sola tabla, crea un receptor (sink) de Cloud Logging con un filtro avanzado hacia Pub/Sub, y luego activa tu consumidor desde ese tema.
Dataform: flujos de trabajo SQL para BigQuery
- Modela grafos de dependencia con
ref(), define tablas/vistas/incrementales y orquesta las compilaciones (builds) por etiquetas o programaciones. Dataform compila SQLX en planes de ejecución ordenados, permitiendo una orquestación basada en metadatos a partir de definiciones declarativas. - Las aserciones garantizan la calidad de los datos. Una aserción es una consulta que debe devolver cero filas para ser aprobada. Ejemplo de aserción: – definitions/assert_precios_no_negativos.sqlx config { type: “assertion” } SELECT 1 FROM ${ref(‘prices_daily’)} WHERE price < 0 LIMIT 1
- Lanzamientos y controles de repositorio: almacena el código en un repositorio, usa ramas y revisiones, y promueve lanzamientos etiquetados a los entornos (p. ej., dev, test, prod) con variables específicas del entorno. Controla los despliegues mediante verificaciones de CI/CD y los resultados de las aserciones.
Dataproc, Dataflow y patrones de almacenamiento
- Para reutilizar Hadoop/Spark con operaciones mínimas, usa Dataproc con el conector de GCS para que los datos persistan más allá del ciclo de vida del clúster y minimizar el costo del disco persistente. Crea clústeres efímeros por trabajo para el aislamiento y el control de costos; orquéstralos con Composer o Workflows.
- Para la ingesta por lotes con filas mal formadas, ejecuta Dataflow para escribir los registros válidos en BigQuery y enrutar los errores de análisis/validación a una tabla de mensajes fallidos (dead-letter) en BigQuery para su inspección.
Observabilidad, alertas y runbooks
Telemetría y alertas
- Enrutar todos los logs de orquestación a Cloud Logging con campos estructurados (pipeline, dag_id, run_id, task_id, partition). Exportar los logs de errores a Monitoring a través de métricas basadas en logs. Alertar sobre:
- Programaciones omitidas o incumplimientos de SLA
- Fallos consecutivos de tareas
- Crecimiento del backlog (p. ej., mensajes sin confirmar [unacked] en Pub/Sub, retraso del sistema [system lag] en Dataflow)
- Fallos en las aserciones de calidad de datos
- Cloud Composer: monitorear la duración de DAGs/tareas, la tasa de éxito, la profundidad de la cola y la salud del planificador (scheduler). Configurar
on_failure_callbackpara notificaciones (paging) y runbooks de remediación. - Cloud Workflows: inspeccionar los logs de ejecución (Execution logs) y las latencias de los pasos; añadir reintentos explícitos y manejadores de errores; emitir logs personalizados con IDs de correlación.
- Notificaciones de cambio en tablas de BigQuery: crear un receptor (sink) de Logging a nivel de proyecto con un filtro avanzado para trabajos de inserción dirigidos a una tabla específica y exportarlo a Pub/Sub; tu herramienta de monitoreo se suscribe al tema para recibir alertas instantáneas sin el ruido de otras tablas.
Diseño de runbooks
- Para cada pipeline, documentar disparadores (triggers), dependencias, SLAs, procedimientos de reversión (rollback)/reintento (retry) y pasos seguros para el rellenado de datos históricos (backfill). Incluir la “repetición de un conjunto de datos fijo” (fixed dataset replay) para Dataflow, cómo drenar un trabajo de streaming, cómo reprocesar particiones fallidas y cómo remediar mensajes de la DLQ.
- Capturar firmas de fallos comunes (p. ej., permiso denegado, cuota excedida, discrepancia de esquema) con árboles de decisión y rutas de escalamiento.
Escenario de un problema práctico
Acme Retail Analytics necesita ingerir entregas diarias de archivos CSV de socios que contienen ocasionalmente filas con formato incorrecto, transformar y cargar los datos válidos a BigQuery, y exponer las filas erróneas para su investigación. También quieren un enriquecimiento basado en eventos para actualizaciones de precios casi en tiempo real y una promoción segura de desarrollo (dev) a producción (prod).
Enfoque:
Almacenamiento y disparadores de eventos
- Crear un bucket de Cloud Storage dedicado con control de versiones de objetos y acceso uniforme a nivel de bucket. Habilitar notificaciones de finalización de objetos (object finalize) a Pub/Sub a través de Eventarc.
- Justificación: La finalización de objetos es un evento fiable para disparar la ingesta posterior; el control de versiones permite reejecuciones y auditorías.
Ingesta por lotes con manejo de cola de mensajes fallidos (dead-letter)
- Usar Cloud Composer para programar un DAG de Airflow diario a las 02:00 con
catchuphabilitado. El DAG lanza un trabajo por lotes de Dataflow que analiza los CSV, valida el esquema y escribe los registros válidos en BigQuery usando tablas de staging deterministas y luego haciendoMERGEen las tablas de destino particionadas. Enrutar los registros con formato incorrecto o fallidos a una tabla de mensajes fallidos (dead-letter) en BigQuery. - Justificación: Dataflow escala el análisis/validación;
MERGEasegura la idempotencia; la captura en una cola de mensajes fallidos permite la inspección sin bloquear el pipeline, coincidiendo con el patrón recomendado para filas con formato incorrecto.
- Usar Cloud Composer para programar un DAG de Airflow diario a las 02:00 con
Enriquecimiento basado en eventos
- Desplegar un trabajo de Cloud Run para realizar un enriquecimiento ligero para actualizaciones de precios incrementales. Dispararlo a través de Cloud Workflows que escucha mensajes de Pub/Sub desde Eventarc cuando llegan pequeños archivos de actualización durante el día.
- Justificación: Los contenedores sin servidor (serverless) con Workflows proporcionan una orquestación de baja latencia y baja sobrecarga operativa para eventos pequeños, mientras se mantienen las transformaciones pesadas en lote.
Controles de fiabilidad
- Configurar reintentos con retroceso exponencial (exponential backoff) para fallos transitorios en los trabajos de Dataflow y Cloud Run, limitando el tiempo total de reintentos al SLA del DAG. Establecer tiempos de espera de ejecución por tarea y callbacks
on_failureen Airflow; en Workflows, establecermax_doublingsymax_retry_duration. - Justificación: El retroceso acotado (bounded backoff) preserva los SLAs y previene reintentos descontrolados.
- Configurar reintentos con retroceso exponencial (exponential backoff) para fallos transitorios en los trabajos de Dataflow y Cloud Run, limitando el tiempo total de reintentos al SLA del DAG. Establecer tiempos de espera de ejecución por tarea y callbacks
Seguridad y mínimo privilegio
- Ejecutar cada componente bajo una cuenta de servicio (service account) dedicada: SA del orquestador de Composer, SA del worker de Dataflow, SA del trabajo de Cloud Run. Otorgar solo los roles necesarios: lectura de GCS (
read) en el bucket de ingesta a Dataflow,dataEditorde BigQuery en los datasets de destino yVieweren los logs. Almacenar los secretos en Secret Manager y hacer referencia a ellos en tiempo de ejecución. - Justificación: Aplica el principio de mínimo privilegio y aísla el radio de impacto (blast radius).
- Ejecutar cada componente bajo una cuenta de servicio (service account) dedicada: SA del orquestador de Composer, SA del worker de Dataflow, SA del trabajo de Cloud Run. Otorgar solo los roles necesarios: lectura de GCS (
Orquestación basada en metadatos
- Mantener una tabla de control en BigQuery que liste las fuentes de los socios, los patrones de archivo y los datasets de destino. En tiempo de ejecución del DAG, Airflow consulta esta tabla y utiliza el mapeo dinámico de tareas para generar tareas por cada socio.
- Justificación: Añadir un socio se convierte en un cambio de datos, no en un cambio de código, reduciendo el riesgo del despliegue.
Observabilidad y alertas
- Emitir logs estructurados con
run_idypartner_id. Crear políticas de alertas para incumplimientos de SLA del DAG, retraso del sistema (system lag) de Dataflow y recuentos no vacíos en la cola de mensajes fallidos. Para las inserciones de BigQuery en la tabla de destino, configurar un receptor (sink) de Cloud Logging con un filtro avanzado para esa tabla hacia un tema de Pub/Sub consumido por la herramienta de monitoreo de Acme. - Justificación: Las alertas detalladas (fine-grained) permiten una clasificación rápida de incidentes (triage) sin ruido.
- Emitir logs estructurados con
CI/CD y promoción
- Gestionar la infraestructura (buckets, Pub/Sub, Eventarc, Composer, Workflows, datasets de BigQuery) en Terraform. Usar Cloud Build para validar la sintaxis de los DAGs de Airflow, ejecutar pruebas unitarias y desplegar a un entorno de desarrollo (dev) de Composer. Promover a los entornos de pruebas (test) y producción (prod) con configuraciones parametrizadas y puertas de aprobación manual después de que las aserciones de Dataform y las pruebas de integración pasen.
- Justificación: Despliegues declarativos y repetibles y promoción segura entre entornos.
Runbook y recuperación
- Documentar los pasos para reprocesar una fecha específica: restaurar el CSV desde el control de versiones de objetos, reejecutar el trabajo de Dataflow para esa partición, hacer
MERGEde los resultados y revisar los registros de la DLQ. Incluir un procedimiento de “repetición de un conjunto de datos fijo” (fixed dataset replay) para aislar errores de transformación si surgen discrepancias. - Justificación: Un diseño idempotente y una recuperación documentada agilizan la remediación de fallos parciales.
- Documentar los pasos para reprocesar una fecha específica: restaurar el CSV desde el control de versiones de objetos, reejecutar el trabajo de Dataflow para esa partición, hacer
← Ingesta · Todos los dominios · Machine Learning →
Practica estas preguntas → · Práctica cronometrada en ExamRoll.io →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
Aprueba tu examen →