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
Pipelines, transformaciones, PCollections, runners:
- Un pipeline de Beam aplica un grafo acíclico dirigido de PTransforms a PCollections (acotadas o no acotadas).
- Los runners (Dataflow, Spark, Flink, Direct) ejecutan el pipeline; Dataflow proporciona autoescalado gestionado, puntos de control (checkpointing) y visibilidad operativa.
- Las transformaciones incluyen operaciones por elemento (ParDo), agrupación y combinación (GroupByKey, Combine), uniones (CoGroupByKey) y E/S (IOs) (PubSubIO, BigQueryIO, FileIO).
Ventanas (Windows):
- Ventanas fijas (Fixed windows): intervalos que no se superponen (p. ej., ventanas de saltos de 1 minuto) para agregados periódicos.
- Ventanas deslizantes (Sliding windows): ventanas que se superponen para métricas continuas (p. ej., ventanas de 5 minutos que se deslizan cada minuto).
- Ventanas de sesión (Session windows): ventanas dinámicas que se cierran tras un intervalo de inactividad, ideales para sesiones de usuario o ráfagas de datos de dispositivos.
- Ventana global (Global window): la vista predeterminada sin ventanas de todo el flujo no acotado; a menudo se combina con disparadores (triggers) para la materialización periódica.
Tiempo de evento vs. tiempo de procesamiento:
- Tiempo de evento: cuándo ocurrió el evento en el origen; permite agregaciones lógicamente consistentes a pesar de las latencias de transporte variables.
- Tiempo de procesamiento: cuándo el evento es observado por el pipeline; útil para disparadores operativos pero no para la correctitud semántica.
Marcas de agua (Watermarks):
- Una marca de agua (watermark) estima la completitud del tiempo de evento (la suposición del runner de que ha visto todos los eventos hasta un tiempo T).
- Las marcas de agua pueden avanzar de forma irregular o detenerse por contrapresión (backpressure) o retrasos en el origen; los datos tardíos son aquellos que llegan con una marca de tiempo < marca de agua.
Disparadores (Triggers) y datos tardíos:
- Predeterminado: disparador AfterWatermark que se activa cuando la marca de agua supera el final de la ventana; con una latencia permitida (allowed lateness) = 0, los datos tardíos se descartan.
- Activaciones tempranas (basadas en tiempo de procesamiento o en conteo) proporcionan resultados preliminares de baja latencia.
- Activaciones tardías permiten correcciones cuando llegan datos tardíos; el modo de acumulación (accumulation mode) gobierna si los paneles (panes) acumulan resultados o descartan la salida anterior.
- Elija la latencia permitida (allowed lateness) según la tolerancia del negocio y el equilibrio entre almacenamiento y cómputo; una mayor latencia aumenta la retención de estado y el costo.
Procesamiento con estado, temporizadores, sesionización y deduplicación:
- Las DoFns con estado (Stateful DoFns) mantienen un estado por clave (p. ej., último evento visto, agregados en curso) y establecen temporizadores para emitir o limpiar el estado.
- La sesionización se expresa de forma natural mediante SessionWindows; para lógica personalizada, use estado por clave y temporizadores de tiempo de procesamiento/evento.
- Deduplicación: use un ID estable por evento y aplique Distinct/Combine por ventana, o un estado por clave (p. ej., un filtro de Bloom o un conjunto con TTL). Equilibre la memoria y los falsos positivos frente a la precisión estricta.
Modos de fallo y contrapartidas:
- Usar ventanas de tiempo de procesamiento para métricas de negocio causa desviaciones durante picos o reintentos; prefiera ventanas de tiempo de evento.
- Ventanas demasiado pequeñas con disparadores tempranos frecuentes causan una emisión excesiva de paneles (panes) y una amplificación de la escritura en el sumidero (sink).
- Una latencia permitida (allowed lateness) ilimitada puede inflar el estado; siempre acote el TTL del estado y configure temporizadores para limpiar las claves inactivas.
Operación de Dataflow para cargas de trabajo de streaming
Dimensionamiento y autoescalado de los trabajadores:
- El autoescalado horizontal añade/elimina trabajadores basándose en el trabajo pendiente (backlog), el retraso de la marca de agua (watermark lag), la CPU y el rendimiento (throughput); establece un
maxWorkerssensato para absorber picos. - Elige tipos de máquina según los cuellos de botella: limitado por CPU (más vCPUs), limitado por memoria (tipos con alta memoria), limitado por red (VMs más grandes reducen la sobrecarga del shuffle).
- Aumenta el disco de arranque para shuffles pesados o receptores basados en archivos. Monitorea el retraso del sistema (system lag) y los segundos de trabajo pendiente (backlog seconds).
- El autoescalado horizontal añade/elimina trabajadores basándose en el trabajo pendiente (backlog), el retraso de la marca de agua (watermark lag), la CPU y el rendimiento (throughput); establece un
Streaming Engine y shuffle:
- Streaming Engine externaliza el estado y el shuffle al backend del servicio, mejorando la elasticidad, reduciendo la presión de memoria en los trabajadores y permitiendo actualizaciones más rápidas.
- Para etapas con mucha carga de lotes o agrupaciones de claves masivas, usa Dataflow Shuffle para descargar la E/S del shuffle de los trabajadores. Ambos reducen las fallas por trabajadores sobrecargados (hot-workers) y la sobrecarga del disco (disk thrash).
Contrapresión, claves calientes y desequilibrio (skew):
- Dataflow gestiona la contrapresión mediante el reequilibrio dinámico del trabajo; no obstante, ajusta el control de flujo en el origen (p. ej., mensajes/bytes pendientes en Pub/Sub) cuando sea aplicable.
- Las claves calientes (p. ej., IDs populares) crean trabajadores rezagados (stragglers). Mitígalo con particionamiento de claves (key#N), preagregación parcial seguida de un re-claveado, o aproximaciones basadas en sketches.
- El desequilibrio (skew) por registros atípicos (cargas útiles enormes) o publicadores con ráfagas de datos puede requerir particiones por publicador, procesamiento por lotes o compresión.
Integración con Pub/Sub:
- Usa temas de Pub/Sub para la ingesta; habilita los atributos de mensaje para los metadatos (p. ej., deviceId, marca de tiempo del evento).
- Realiza la ingesta con PubSubIO; extrae las marcas de tiempo del evento de los atributos o de la carga útil; si no, recurre al tiempo de publicación.
- Las claves de ordenamiento (ordering keys) proporcionan orden por clave; Dataflow aún necesita un comportamiento idempotente en los sistemas de destino debido a la entrega de tipo “al menos una vez” (at-least-once).
Patrones de streaming a BigQuery:
- Prefiere BigQueryIO con la Storage Write API para un alto rendimiento, baja latencia y semántica “exactamente una vez” (exactly-once) dentro de un stream mediante desplazamientos (offsets) de stream y reintentos automáticos.
- Para canalizaciones simples de baja tasa, las inserciones de streaming son aceptables; establece un
insertIdpara deduplicar los reintentos del cliente. - Las consultas sobre los búferes de streaming son eventualmente consistentes; para análisis críticos en el tiempo, consulta después de un retardo del búfer (p. ej., espera ~2x la latencia de disponibilidad observada), o materializa mediante ventanas de microlotes y el modo confirmado (committed mode) de la Storage Write API.
Efectos de “exactamente una vez”, idempotencia, reproducción y receptores:
- Beam garantiza un procesamiento de “al menos una vez”; el “exactamente una vez” debe lograrse en el receptor usando escrituras idempotentes, transacciones o claves de deduplicación.
- BigQuery: usa los streams por defecto o los streams confirmados de la Storage Write API para obtener semántica “exactamente una vez” dentro de un stream; con inserciones de streaming, establece un
insertIdestable. - Archivos: escribe archivos temporales con nombres únicos, finalízalos al completarse la ventana y asegura renombramientos atómicos; evita sobrescribir para prevenir duplicados parciales.
- Bases de datos externas: usa operaciones “upsert” (actualizar o insertar) basadas en un ID estable o implementa ventanas de deduplicación.
- Diseña para la reproducción: mantén transformaciones deterministas; asegúrate de que los receptores dedupliquen en los reintentos.
Manejo de mensajes fallidos (dead-letter), enrutamiento de errores y observabilidad:
- Envuelve el análisis (parsing) o enriquecimiento riesgoso en un bloque try/catch dentro de un ParDo y emite los fallos a una PCollection de mensajes fallidos (dead-letter) mediante un TupleTag; incluye la carga útil, el código de error y el contexto.
- Enruta las colas de mensajes fallidos (DLQs) a BigQuery o Cloud Storage para su análisis; considera un tema de Pub/Sub separado para el reprocesamiento.
- Observabilidad: usa las métricas de trabajo de Dataflow (retraso de la marca de agua, retraso del sistema, rendimiento), contadores personalizados, métricas de distribución y registros por paso en Cloud Logging. Crea alertas sobre el retraso y las tasas de error en Cloud Monitoring. Usa Error Reporting para agregar las excepciones.
Patrones de ajuste de rendimiento:
- Lee eficientemente: para orígenes de BigQuery, prefiere la Storage Read API o lecturas basadas en consultas que seleccionen solo los campos y filtros necesarios.
- Optimización de combinación (Combine lifting): usa combinadores (combiners) para reducir el volumen del shuffle antes de un GroupByKey.
- Entradas laterales (Side inputs): almacena en caché datos de referencia pequeños en memoria; vigila el factor de distribución (fanout) y la cadencia de actualización.
- Serialización: usa esquemas compactos (Avro/Proto) y evita el análisis excesivo de JSON en las rutas críticas (hot paths).
Estrategias de despliegue, plantillas y actualización
Flex Templates:
- Empaquetan las canalizaciones en plantillas contenerizadas y parametrizadas para despliegues reproducibles. Las Flex Templates soportan dependencias personalizadas, imágenes de GPU y aislamiento de entorno.
- Externalizar los parámetros de tiempo de ejecución (p. ej., suscripción de entrada, tabla de salida, receptor de mensajes fallidos o dead-letter sink, maxWorkers) para permitir despliegues específicos para cada entorno.
Actualizaciones y compatibilidad de canalizaciones:
- Dataflow soporta la actualización in-situ para muchas canalizaciones de streaming si los nombres de las transformaciones, las especificaciones de estado y los tipos de salida permanecen compatibles. Utilice nombres de PTransform estables.
- Para cambios incompatibles en el grafo o el estado, realice una transición controlada: inicie el nuevo trabajo, luego drene el trabajo antiguo para finalizar el trabajo en curso y dejar de leer nuevos elementos.
Drenaje e instantáneas (snapshots):
- El drenaje completa el procesamiento de forma controlada, escribe la salida restante y termina; coordine con la retención o las instantáneas de Pub/Sub para evitar brechas de datos.
- Para garantizar la continuidad, puede crear una instantánea de Pub/Sub, iniciar la nueva canalización posicionándose en la instantánea o en una marca de tiempo apropiada, verificar la salida y luego drenar el trabajo antiguo.
Ejemplos de configuración:
- Ejemplo de ventaneo con disparadores tempranos/tardíos y acumulación:
undefined
- Ejemplo de BigQueryIO con la Storage Write API:
undefined
- Errores comunes:
- Escribir en receptores basados en archivos en modo streaming sin escrituras en ventana puede bloquear la finalización; habilite las escrituras en ventana y los disparadores.
- Crecimiento ilimitado: olvidar limitar el estado o la latencia permitida puede causar fugas de memoria y fallos de escalado.
- Marcas de tiempo ausentes: no asignar marcas de tiempo de evento hace que la canalización utilice por defecto el tiempo de procesamiento y pierda correctitud bajo retardos variables.
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:
Ingesta y semántica de tiempo
- Crear un tema regional de Pub/Sub y publicadores por región con los atributos
deviceIdyeventTs(RFC3339). Habilitar claves de ordenamiento pordeviceIdcuando 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.
- Crear un tema regional de Pub/Sub y publicadores por región con los atributos
Canalización de streaming de Dataflow con ventanas de tiempo de evento
- Leer desde una suscripción dedicada a través de PubSubIO, extrayendo
eventTscomo la marca de tiempo de Beam, y recurriendo apublishTimesi 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.
- Leer desde una suscripción dedicada a través de PubSubIO, extrayendo
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
deviceIdy 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.
Agregación y mitigación de hot keys
- Agrupar por
deviceIdy calcular avg/min/max por minuto con CombineFns. Para las métricas regionales top-N, fragmentar porregion#Npara 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).
- Agrupar por
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
insertIdestable basado endeviceId+eventTspara 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.
- Escribir los eventos validados en crudo y los agregados por minuto en BigQuery usando BigQueryIO con la Storage Write API. Establecer un
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.
Operaciones: autoescalado y Streaming Engine
- Habilitar Streaming Engine; establecer
maxWorkersbasá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.
- Habilitar Streaming Engine; establecer
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,
maxWorkersy 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.
- Empaquetar la canalización como una Flex Template con parámetros: suscripción de entrada, tablas de salida, tabla de DLQ,
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 →