Google PDE: Mensajería, Ingesta de Eventos y Servicios en Tiempo Real — 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
Los servicios de mensajería, ingesta de eventos y en tiempo real en Google Cloud se centran en Cloud Pub/Sub y Eventarc para un transporte desacoplado y duradero; Dataflow para el procesamiento de flujos con estado; y receptores como BigQuery, Cloud Storage y bases de datos operativas. Diseñar para una entrega de tipo “al menos una vez”, un consumo idempotente y una buena observabilidad garantiza sistemas resilientes que escalan elásticamente mientras mantienen la corrección ante fallos, contrapresión y evolución de esquemas.
Mensajería principal con Pub/Sub
- Temas y suscripciones
- Los publicadores envían mensajes a un tema; los suscriptores se conectan a través de suscripciones (múltiples suscriptores pueden consumir los mismos mensajes de forma independiente).
- Tipos de suscripción:
- Pull: los clientes solicitan explícitamente los mensajes; utiliza streaming pull para obtener el máximo rendimiento y menos viajes de ida y vuelta.
- Push: Pub/Sub entrega a través de HTTPS; tu punto de conexión (endpoint) debe devolver un código 2xx para confirmar la recepción.
- Exportar a BigQuery: una suscripción de BigQuery entrega mensajes a una tabla de BigQuery sin necesidad de código; es la mejor opción cuando las cargas útiles coinciden con el esquema declarado y se requiere una ingesta de baja latencia para el análisis.
- Claves de ordenación
- Habilita el orden de los mensajes en el tema y la suscripción para recibir entregas en orden por cada clave de ordenación. El rendimiento por clave se serializa: un mensaje en tránsito por clave puede bloquear los siguientes; utiliza muchas claves (por ejemplo, hash(device_id)) para escalar.
- Distribución (fan-out) y reproducción (replay)
- Crea suscripciones separadas para diferentes consumidores para aislar las cargas de trabajo y la retención.
- Usa ‘seek’ o ‘snapshot’ para reproducir desde una marca de tiempo o una instantánea para recuperación y rellenos de datos (backfills).
Compensaciones:
- La ordenación reduce el paralelismo y el rendimiento por clave; desactiva la ordenación a menos que sea estrictamente necesario.
- El modo Push simplifica el código del cliente, pero introduce problemas de escalado del punto de conexión HTTP, seguridad y retroceso exponencial (backoff); el modo Pull ofrece más control y estabilidad con un alto rendimiento.
Semántica de entrega, confirmación, retención y colas de mensajes fallidos (Dead Lettering)
- Confirmación (acknowledgment) y plazos
- Entrega de tipo “al menos una vez”: pueden producirse duplicados.
- Cada entrega tiene un plazo de confirmación (ack deadline) (10 segundos por defecto). Extiéndelo (ModifyAckDeadline) mientras procesas trabajos de larga duración; no confirmar la recepción antes del plazo es la causa más común de entregas Push duplicadas.
- Un Nack (confirmación negativa) o el vencimiento del plazo hacen que el mensaje sea elegible para una nueva entrega.
- Retención
- Los mensajes no confirmados se retienen durante el plazo de confirmación de la suscripción y se reintentan; los mensajes confirmados pueden retenerse hasta la duración de retención de mensajes del tema para su reproducción. Configura la retención para cubrir tu tiempo máximo de interrupción más el tiempo de recuperación.
- Reintentos
- Pull: la reentrega ocurre después de que expire el plazo de confirmación; controla la concurrencia con límites de control de flujo.
- Push: retroceso exponencial (exponential backoff); solo un código HTTP 2xx se considera éxito. Los códigos 3xx/4xx/5xx desencadenan reintentos. Implementa manejadores idempotentes para tolerar repeticiones.
- Temas de mensajes fallidos (DLT)
- Configura un tema DL y un número máximo de intentos de entrega por suscripción para poner en cuarentena los mensajes “venenosos” (poison messages).
- Supervisa el volumen de la DLQ (cola de mensajes fallidos); crea flujos de trabajo de triaje y vuelve a publicar en el tema principal después de la corrección.
Ejemplo:
undefined
Resumen de la semántica de entrega:
- Pub/Sub: entrega de tipo “al menos una vez”, ordenación de mejor esfuerzo dentro de una clave de ordenación si está habilitada.
- Receptores: las API de inserción de BigQuery proporcionan mitigación de duplicados (insertId o desplazamientos de flujo de la Storage Write API), pero aun así, diseña los consumidores y escritores para que sean idempotentes.
Esquemas, compatibilidad y validación
- Esquemas de Pub/Sub
- Soporte nativo para Avro y Protocol Buffers con esquemas almacenados de forma centralizada.
- Configuración del esquema a nivel de tema: codificación (Avro o Protobuf) y aplicación (ninguna, solo validar o requerir).
- El productor publica cargas útiles codificadas; Pub/Sub valida contra el esquema actual cuando la aplicación está habilitada.
- Evolución y compatibilidad
- Utiliza cambios retrocompatibles (añadir campos opcionales, añadir campos con valores predeterminados en Avro, no reutilizar nunca las etiquetas en Protobuf, evitar eliminar o renombrar campos).
- Versiona los esquemas explícitamente. Para cambios disruptivos (breaking changes), publica dualmente en temas v1 y v2, o añade un campo de versión y enruta en consecuencia.
- Contratos productor-consumidor
- Los consumidores deben ignorar los campos desconocidos y asignar valores predeterminados a los que falten.
- Prueba la compatibilidad del esquema con todos los consumidores antes de promoverlo a producción; valida en suscripciones de preproducción (staging) con la misma aplicación de esquema que en producción.
Ejemplo corto de Avro (extracto):
undefined
Integración basada en eventos, Eventarc e interoperabilidad con Kafka
- Eventarc y CloudEvents
- Eventarc enruta eventos desde servicios de Google Cloud (y fuentes personalizadas a través de Pub/Sub) hacia Cloud Run, GKE o Workflows utilizando la especificación de CloudEvents. Atributos como
type,sourceysubjectpermiten un filtrado de grano fino y auditabilidad. - Usa filtros de atributos para minimizar el fan-out y reducir la carga en los sistemas de destino (downstream).
- La entrega es at-least-once (al menos una vez); haz que los manejadores (handlers) sean idempotentes y sin estado (stateless) siempre que sea posible.
- Eventarc enruta eventos desde servicios de Google Cloud (y fuentes personalizadas a través de Pub/Sub) hacia Cloud Run, GKE o Workflows utilizando la especificación de CloudEvents. Atributos como
- Ejemplo de un trigger de Eventarc:
gcloud eventarc triggers create gcs-finalize-to-run
–destination-run-service=ingestor
–event-filters=“type=google.cloud.storage.object.v1.finalized”
–event-filters=“bucket=my-data-bucket”
–service-account=eventarc-sa@PROJECT_ID.iam.gserviceaccount.com - Interoperabilidad con Kafka y migración gestionada
- Las plantillas de Dataflow conectan Kafka <-> Pub/Sub para una migración por fases. Replica los topics conservando las claves; transfiere primero los consumidores, luego los productores, o utiliza escritura dual (dual-write) durante la transición.
- Pub/Sub Lite ofrece streaming particionado y con capacidad aprovisionada, con enrutamiento basado en claves y un costo menor; es regional/zonal y es adecuado para cargas de trabajo similares a Kafka donde la capacidad predecible y el orden por partición son las principales prioridades.
- Consideraciones para la migración:
- Ordenamiento: mapea las claves de Kafka a las claves de ordenamiento (ordering keys) de Pub/Sub o a las particiones de Lite.
- Offsets: transporta los offsets como atributos del mensaje para diagnósticos; los consumidores no pueden depender de los offsets de Kafka después de la migración.
- Entrega: acepta la entrega at-least-once; implementa la idempotencia en los sistemas de destino (downstream).
- Esquemas: migra las definiciones de Confluent Schema Registry a los esquemas de Pub/Sub o estandariza en Protobuf/Avro con reglas de evolución compatibles.
Patrones de ingesta en streaming, rendimiento, escalado, seguridad y operaciones
- Patrones de ingesta en tiempo real
- Pub/Sub -> Dataflow -> BigQuery: usa el sink de la API de escritura de almacenamiento (Storage Write API) de BigQuery para un alto rendimiento e idempotencia con offsets de stream; redirige los fallos a una tabla de mensajes fallidos (dead-letter) para su inspección.
- Pub/Sub -> Dataflow -> Cloud Storage: archiva eventos sin procesar para su reprocesamiento; utiliza escrituras en ventanas y comprimidas para equilibrar el costo y la latencia.
- Pub/Sub -> almacenes operacionales: escribe en Bigtable para búsquedas de baja latencia, en Spanner para transacciones fuertemente consistentes, o en Cloud SQL/Firestore según las necesidades de la carga de trabajo. Asegura upserts idempotentes usando un ID de evento único como clave.
- Entrega de tipo at-least-once, prevención de duplicados e idempotencia
- Incluye un
event_idy unevent_timeúnicos en cada mensaje; impón el uso de UUID en el lado del productor. - Deduplicación en el streaming de BigQuery: establece el
insertIdo usa la Storage Write API con streams ordenados; aun así, protege las consultas con lógica de deduplicación. - Ejemplo de deduplicación en tiempo de consulta:
- Incluye un
undefined
- Para los endpoints de tipo push, devuelve un código 2xx solo después de un procesamiento exitoso; de lo contrario, se debe esperar una reentrega.
- Rendimiento de mensajes, cuotas y escalado
- Publicadores (Publishers): agrupa los mensajes en lotes y reutiliza las conexiones; paraleliza entre múltiples clientes. Usa muchas claves de ordenamiento (ordering keys) para escalar cargas de trabajo ordenadas.
- Suscriptores (Subscribers): prefiere el modo streaming pull con control de flujo (max outstanding bytes/messages). Dimensiona los plazos de confirmación (ack deadlines) según el tiempo de procesamiento y extiéndelos cuando sea necesario.
- Monitoriza y solicita aumentos de cuota para el rendimiento de publicación y suscripción a medida que los volúmenes crecen; diseña con un margen de capacidad (por ejemplo, 2x el pico esperado) para absorber ráfagas.
- Consistencia y disponibilidad
- El streaming de BigQuery es eventualmente consistente para la visibilidad de las consultas; para consultas interactivas que deben incluir filas recién insertadas por streaming, espera un tiempo basado en la latencia observada (por ejemplo, 2 veces el retardo de disponibilidad P50) o diseña usando agregaciones alineadas con marcas de agua (watermarks) en Dataflow y consulta los resultados materializados.
- Seguridad
- IAM: concede roles de privilegio mínimo (
pubsub.publishera los productores en el tema;pubsub.subscribera los consumidores en la suscripción). Usa cuentas de servicio dedicadas para cada carga de trabajo. - Autenticación push: configura las suscripciones de tipo push para adjuntar tokens OIDC de una cuenta de servicio; exige la validación de la audiencia (audience) en el endpoint. Prefiere los endpoints privados de Cloud Run para obtener autenticación y TLS integrados.
- Cifrado: Pub/Sub cifra los datos en tránsito y en reposo; usa CMEK en los temas para claves gestionadas por el cliente. Aplica VPC Service Controls para reducir el riesgo de exfiltración de datos. Usa cifrado del lado del cliente para campos sensibles de la carga útil (payload) si es necesario.
- IAM: concede roles de privilegio mínimo (
- Diagnóstico operativo de retraso (lag), reentregas y fallos del suscriptor
- Monitoriza con Cloud Monitoring:
subscription/num_undelivered_messagesyoldest_unacked_message_agepara el trabajo pendiente (backlog).expired_ack_deadline_countpara detectar confirmaciones (acks) omitidas que causan duplicados.publish_request_countypull_request_countpara el rendimiento.
- Investiga eventos que no aparecen en el dashboard reproduciendo un conjunto de datos conocido a través del pipeline y comparando las salidas de cada etapa para aislar la transformación o el sink defectuoso.
- Para el streaming de Dataflow:
- Usa el autoescalado con un
maxWorkersapropiado para absorber la carga de muchas fuentes. - Drena los pipelines (drain) para actualizaciones incompatibles para permitir que el trabajo en curso se complete y evitar la pérdida de datos.
- Usa el autoescalado con un
- Para las notificaciones de inserción de BigQuery, redirige las entradas de auditoría de Cloud Logging a través de un sink filtrado a tablas específicas hacia un tema de Pub/Sub para generar alertas.
- Monitoriza con Cloud Monitoring:
Escenario de un problema práctico
Contoso Freight necesita una plataforma global de eventos en tiempo real para ingerir 10,000 mensajes de telemetría de IoT por minuto desde camiones, enriquecer eventos, potenciar análisis interactivos y activar flujos de trabajo cuando socios externos dejen archivos. Algunos CSV de los socios contienen filas con formato incorrecto, y el equipo de análisis debe inspeccionar los errores sin bloquear el stream.
- Crear la capa principal de mensajería y esquemas
- Acción: Define un esquema Avro para la telemetría y adjúntalo a un tema de Pub/Sub llamado
telemetrycon la validación de esquema (schema enforcement) configurada como obligatoria. Habilita el ordenamiento de mensajes y publica conordering_key = hash(device_id). - Justificación: La validación de esquema a nivel de tema rechaza los eventos con formato incorrecto de forma temprana. El ordenamiento por dispositivo permite el procesamiento ordenado cuando es necesario, mientras que el hashing distribuye las claves para preservar el rendimiento.
- Aprovisionar suscripciones con aislamiento y gestión de mensajes fallidos (dead-lettering)
- Acción: Crea una suscripción de tipo pull
telemetry-stream-subpara Dataflow con un tema de mensajes fallidos (dead-letter)telemetry-dltymax_delivery_attempts=10. Añade una suscripción de BigQuerytelemetry-raw-bqpara depositar los eventos sin procesar en una tabla particionada por tiempo para linaje y reproducción. - Justificación: La DLQ (cola de mensajes fallidos) aísla los mensajes problemáticos (poison messages) para su investigación. Una suscripción de BigQuery separada proporciona una ruta de exportación de baja sobrecarga operativa para la preservación de eventos sin procesar, independientemente del pipeline de procesamiento.
- Construir un pipeline de streaming de Dataflow para enriquecimiento y sinks
- Acción: Ingiere desde
telemetry-stream-subusando streaming pull con control de flujo. Valida contra el esquema, enriquece con datos de referencia y calcula agregados en ventanas. Escribe en BigQuery usando la Storage Write API con un stream con nombre yinsertId = event_id; escribe copias de seguridad de los datos sin procesar en Cloud Storage cada hora; redirige los registros incorrectos/fallidos a una tabla de mensajes fallidos en BigQuery. - Justificación: La Storage Write API ofrece escrituras de alto rendimiento y baja latencia con idempotencia a través de
insertId/offsets del stream. Una tabla de mensajes fallidos permite la inspección sin bloquear el stream, y los archivos en Cloud Storage permiten la reproducción.
- Manejar duplicados y consistencia eventual en los análisis
- Acción: Para consultas interactivas que deben excluir duplicados, publica
event_idyevent_timeen cada registro y usa una vista de deduplicación:
undefined
Introduce un breve retardo en la consulta basado en la disponibilidad observada del streaming de BigQuery (por ejemplo, el doble de la latencia mediana).
- Justificación: La entrega de tipo at-least-once (al menos una vez) requiere escrituras idempotentes y deduplicación en tiempo de consulta. Esperar reduce la omisión de filas en tránsito, dada la latencia de visibilidad del streaming.
- Integrar la entrega de archivos de socios con Eventarc
- Acción: Configura Eventarc para enrutar los eventos
object.finalizedde Cloud Storage para el bucketpartner-dropsa un servicio de Cloud Run que lanza un trabajo por lotes (batch) de Dataflow para cargar los CSV en BigQuery, enviando los errores de análisis a una tabla de mensajes fallidos. - Justificación: Eventarc proporciona orquestación basada en eventos con filtrado de CloudEvents por bucket y prefijo de objeto. Un trabajo por lotes de Dataflow separa las filas con formato incorrecto para su análisis mientras carga los datos correctos rápidamente.
- Asegurar la plataforma
- Acción: Usa cuentas de servicio distintas: los productores obtienen el rol
pubsub.publisheren el tematelemetry; la SA del worker de Dataflow obtienepubsub.subscriberentelemetry-stream-suby acceso de escritura a los datasets de BigQuery y a Cloud Storage de destino; el trigger de Eventarc usa una SA dedicada con el rol de invocador (invoker) en Cloud Run. Habilita CMEK en el tematelemetryy en los datasets de BigQuery. Configura los endpoints de tipo push, si los hay, con OIDC y validación de audiencia (audience). - Justificación: El principio de privilegio mínimo en IAM y el uso de CMEK cumplen con los requisitos de seguridad y cumplimiento; la entrega autenticada previene la suplantación de identidad (spoofing).
- Operar y escalar de forma fiable
- Acción: Configura el autoescalado de Dataflow con un
maxWorkersgeneroso para absorber picos. Monitorizasubscription/oldest_unacked_message_ageyexpired_ack_deadline_count; genera alertas cuando se superen los umbrales. Para cambios en el pipeline que rompan la compatibilidad, despliega con la opción de drenado (drain) para evitar la pérdida de mensajes. Si el retraso (lag) aumenta, incrementa el paralelismo de los suscriptores y extiende los plazos de confirmación (ack deadlines) proporcionalmente al tiempo de procesamiento. - Justificación: La monitorización proactiva detecta retrasos y reentregas de forma temprana. El autoescalado y los plazos de confirmación ajustados previenen tormentas de duplicados. El drenado preserva los mensajes en tránsito durante las actualizaciones.
Este diseño proporciona una ingesta en tiempo real resiliente, segura y observable con integración por lotes basada en eventos, soporta la tolerancia a duplicados y la evolución de esquemas, y ofrece análisis rápidos mientras aísla los datos incorrectos para una corrección específica.
← Procesamiento de Flujos con Dataflow y Apache Beam · Todos los dominios · Spark →
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 →