Google PDE: Procesamiento de Flujos con Dataflow y Apache Beam — 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.

Resumen general

El procesamiento de flujos (stream processing) en Google Cloud se centra en el modelo de programación unificado de Apache Beam, ejecutado por el runner de Dataflow. Beam proporciona una abstracción lógica —pipelines de transformaciones sobre PCollections— que desacopla tu código de los detalles de ejecución como el paralelismo, el autoescalado y la tolerancia a fallos. En el streaming, la correctitud depende de la semántica del tiempo (tiempo de evento vs. tiempo de procesamiento), el ventaneo (fijo, deslizante, de sesión, global), las marcas de agua (watermarks), los disparadores (triggers) y el manejo de datos tardíos. La excelencia operativa en Dataflow requiere el dimensionamiento correcto de los workers, la política de autoescalado adecuada, el motor de streaming, las elecciones de shuffle, un diseño de sumidero (sink) idempotente, el manejo de colas de mensajes fallidos (dead-letter) y una observabilidad robusta.

Modelo de Apache Beam y semántica del tiempo

Modos de fallo y contrapartidas:

Operación de Dataflow para cargas de trabajo de streaming

Estrategias de despliegue, plantillas y actualización

undefined

undefined

Escenario de un Problema Práctico

NovaTrack Inc. ingiere telemetría global de IoT de 50,000 sensores de temperatura y debe entregar agregados a nivel de minuto, persistir los datos en crudo y presentar un panel de control en tiempo real. Se esperan mensajes malformados ocasionales y entrega fuera de orden. La solución debe autoescalar, exponer los registros erróneos para su inspección y soportar actualizaciones sin tiempo de inactividad.

Enfoque:

  1. Ingesta y semántica de tiempo

    • Crear un tema regional de Pub/Sub y publicadores por región con los atributos deviceId y eventTs (RFC3339). Habilitar claves de ordenamiento por deviceId cuando sea factible.
    • Justificación: Pub/Sub proporciona una entrada de datos duradera y elástica con entrega at-least-once (al menos una vez). Adjuntar las marcas de tiempo del evento en el borde preserva el tiempo real del evento; el ordenamiento por dispositivo reduce el reordenamiento dentro del mismo dispositivo sin cuellos de botella centrales.
  2. Canalización de streaming de Dataflow con ventanas de tiempo de evento

    • Leer desde una suscripción dedicada a través de PubSubIO, extrayendo eventTs como la marca de tiempo de Beam, y recurriendo a publishTime si falta.
    • Aplicar FixedWindows de 1 minuto con un disparador temprano a los 30 segundos y activaciones tardías por cada elemento tardío; establecer una latencia permitida de 10 minutos y paneles acumulativos.
    • Justificación: Las ventanas de tiempo de evento aseguran agregados por minuto precisos; las activaciones tempranas alimentan el panel de control con una frescura por debajo del minuto; las activaciones tardías corrigen los agregados a medida que llegan datos retrasados. El límite de latencia acota el tamaño del estado y el costo.
  3. Validación, enriquecimiento y enrutamiento a cola de mensajes fallidos (dead-letter)

    • Implementar un ParDo que analice el JSON, valide el esquema y los rangos, y enriquezca con datos de referencia estáticos pequeños a través de una entrada lateral (side input) cargada desde BigQuery al inicio del trabajo.
    • Usar TupleTags para emitir registros válidos a la salida principal y los fallos a una PCollection de mensajes fallidos (dead-letter) que contenga la carga útil, el error, el deviceId y la marca de tiempo del análisis; escribir la DLQ a una tabla particionada de BigQuery.
    • Justificación: Las entradas laterales mantienen los datos de referencia en memoria para una baja latencia. La captura de mensajes fallidos permite la inspección y el reprocesamiento dirigido de filas erróneas sin bloquear el flujo principal.
  4. Agregación y mitigación de hot keys

    • Agrupar por deviceId y calcular avg/min/max por minuto con CombineFns. Para las métricas regionales top-N, fragmentar por region#N para evitar hot keys, y luego volver a agregar.
    • Justificación: Los combinadores minimizan el volumen de shuffle y el costo; la fragmentación de claves previene cuellos de botella de una sola clave durante la agregación regional (fan-in).
  5. Receptores y efectos exactly-once

    • Escribir los eventos validados en crudo y los agregados por minuto en BigQuery usando BigQueryIO con la Storage Write API. Establecer un insertId estable basado en deviceId + eventTs para la idempotencia en cualquier reintento personalizado.
    • Justificación: La Storage Write API proporciona una ingesta de alto rendimiento y baja latencia con semántica exactly-once dentro de un flujo. Los ID estables aseguran la deduplicación downstream si ocurren repeticiones.
  6. Estrategia de consistencia del panel de control

    • El panel de control consulta las tablas de agregados particionadas con una ventana de retrospectiva de 2 minutos en relación con la marca de agua (watermark) o un retardo fijo de 2 veces la latencia de disponibilidad observada para los datos en streaming.
    • Justificación: La visibilidad del streaming en BigQuery es eventualmente consistente; diferir las lecturas ligeramente evita la pérdida de filas en tránsito mientras se mantiene un comportamiento casi en tiempo real.
  7. Operaciones: autoescalado y Streaming Engine

    • Habilitar Streaming Engine; establecer maxWorkers basándose en el pico esperado (p. ej., 3x el promedio), seleccionar un tipo de máquina dimensionado para el análisis y cifrado intensivos en CPU, y aumentar el disco de arranque para acomodar el shuffle transitorio.
    • Monitorear el retraso de la marca de agua (watermark lag), los segundos de trabajo pendiente (backlog), la CPU y el rendimiento por paso; alertar sobre retrasos sostenidos y picos en la tasa de DLQ.
    • Justificación: Streaming Engine externaliza el estado/shuffle para mayor elasticidad y actualizaciones más simples; el dimensionamiento correcto y la monitorización previenen incumplimientos silenciosos de SLO.
  8. Despliegue y actualizaciones con Flex Templates

    • Empaquetar la canalización como una Flex Template con parámetros: suscripción de entrada, tablas de salida, tabla de DLQ, maxWorkers y región. Para un cambio incompatible, iniciar la nueva canalización apuntando al mismo tema con una nueva suscripción, verificar las salidas y luego drenar el trabajo antiguo. Opcionalmente, crear una instantánea de Pub/Sub y posicionar la nueva suscripción en la instantánea para garantizar que no haya brechas de datos.
    • Justificación: Las Flex Templates permiten despliegues repetibles y parametrizados. Una transición azul/verde (blue/green) verificada con drenaje logra cero pérdida de datos y un tiempo de inactividad mínimo.
  9. Reprocesamiento y rellenos de datos (backfills) en lote

    • Almacenar archivos Avro comprimidos de eventos en crudo en Cloud Storage a través de una salida lateral; ejecutar una canalización de Dataflow en lote para rellenar o reprocesar datos en BigQuery cuando los modelos o esquemas cambien.
    • Justificación: Los archivos de datos en crudo duraderos soportan la reproducibilidad y la evolución del esquema sin impactar la ruta crítica (hot path).

Este diseño produce agregados correctos y de baja latencia con un costo acotado, un claro aislamiento de errores, una fuerte observabilidad y rutas de actualización seguras, mientras maneja datos fuera de orden y tardíos a escala global.


Analítica de BigQuery e Ingeniería de Almacenes de Datos · Todos los dominios · Mensajería

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 →

Explorar Google →

Related guides

Acceso todo en uno

Una suscripción. Todos los exámenes.

Cada plan desbloquea la búsqueda ilimitada de respuestas, pruebas de práctica, explicaciones de AI y la biblioteca completa de recursos, en más de 20 idiomas.

Mensual
24.87
Just €0.83/day
Todo incluido:
  • Búsqueda ilimitada de respuestas
  • Pruebas de práctica ilimitadas
  • Explicaciones con tecnología AI
  • Biblioteca completa de recursos
  • Más de 20 idiomas
  • Actualizaciones semanales de contenido
  • Recompensas y referencias
  • Soporte prioritario
Iniciar prueba gratuita

No se requiere tarjeta de crédito*

Mejor valor
12 meses
179.87
Just €0.49/daySave 40%
Todo incluido:
  • Búsqueda ilimitada de respuestas
  • Pruebas de práctica ilimitadas
  • Explicaciones con tecnología AI
  • Biblioteca completa de recursos
  • Más de 20 idiomas
  • Actualizaciones semanales de contenido
  • Recompensas y referencias
  • Soporte prioritario
Iniciar prueba gratuita

No se requiere tarjeta de crédito*

✓ Plan gratuito incluido · ✓ Cancela en cualquier momento · ✓ Todos los planes desbloquean el producto completo