Google PDE: Ingesta, Integración y Migración de Datos — 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 ingesta, integración y migración de datos en Google Cloud abarcan patrones repetibles, servicios gestionados y controles operativos que convierten diversos sistemas de origen en conjuntos de datos fiables y consultables. Los diseños eficaces separan el transporte de la transformación, desacoplan productores y consumidores, y favorecen las canalizaciones idempotentes y con puntos de control (checkpointed) con un linaje y una verificación claros. Esta sección cubre los patrones de ingesta, los servicios de Google Cloud para el movimiento y CDC, los controles de esquema y calidad de datos, la conectividad y la integración híbrida, y las estrategias de transición (cutover), destacando en todo momento las contrapartidas de diseño y los modos de fallo.
Patrones y cargas de trabajo de ingesta
- Ingesta por lotes: Extracciones periódicas o entregas de archivos a intervalos definidos. Adecuado para un coste predecible y rellenos de datos históricos (backfills). Modo de fallo: lotes grandes y poco frecuentes causan picos de recursos, largas ventanas de recuperación y el incumplimiento de los SLAs. Mitigación: dimensionar correctamente las ventanas de los lotes, fragmentar (shard) por tiempo o clave, y usar paralelismo.
- Carga masiva: Cargas únicas o a gran escala (p. ej., relleno histórico inicial). Prefiera formatos columnares o autodescriptivos (Parquet, Avro) y cargue directamente en el almacenamiento analítico (BigQuery) o en un área de preparación (staging) en Cloud Storage. Contrapartida: consultar tablas externas evita los pasos de carga, pero traslada el coste al escaneo en tiempo de consulta.
- Carga incremental: Cargas periódicas de deltas mediante marcas de tiempo (timestamps) o marcas de agua alta (high-water marks). Requiere una deduplicación robusta y operaciones de inserción/actualización (upserts) idempotentes. Modo de fallo: desfase de reloj o registros que llegan con retraso. Use marcas de tiempo de confirmación (commit) del lado del servidor y marcas de agua (watermarking).
- Captura de datos de cambios (CDC): Replicación continua de inserciones, actualizaciones y eliminaciones desde bases de datos operacionales. Ideal para análisis casi en tiempo real y migraciones con bajo tiempo de inactividad. Contrapartidas:
- Ordenación: La mayoría de las herramientas de CDC preservan el orden dentro de las transacciones y, típicamente, dentro de un fragmento (shard), pero no garantizan un orden global entre fragmentos. Use marcas de tiempo de confirmación de la transacción y claves primarias para reconstruir la secuencia.
- Semántica de entrega: La entrega at-least-once (al menos una vez) es lo habitual; construya receptores (sinks) idempotentes o deduplique usando identificadores de cambio únicos.
- Snapshot + CDC: Comience con una instantánea (snapshot) consistente, luego aplique los cambios desde una secuencia de registro precisa para alcanzar la paridad sin tiempo de inactividad.
Fuentes relacionales, SaaS, on-premises y de archivos:
- Fuentes relacionales: Use CDC nativo o columnas de marca de tiempo. Para cargas masivas, exporte a Avro/Parquet y almacene en una zona intermedia (stage) en Cloud Storage.
- Fuentes SaaS: Prefiera las APIs del proveedor con tokens incrementales; integre a través de conectores gestionados (p. ej., en Data Fusion). Controle el flujo para no superar los límites de peticiones (rate limits) y gestione la deriva del esquema (schema drift).
- Fuentes on-prem: Elija entre transferencia basada en agentes, VPN/Interconnect + Private Google Access, o siembra de datos offline con Transfer Appliance.
- Ingesta de archivos: Para muchos archivos pequeños, agrúpelos (p. ej., con tar) para reducir la sobrecarga de RPC. Use gsutil -m o clientes paralelizados; componga o transforme en archivos columnares más grandes para el análisis.
Servicios de Google Cloud para ingesta, integración y migración
- Datastream (CDC sin servidor): Captura cambios de MySQL, PostgreSQL y Oracle en Cloud Storage, BigQuery (a través de plantillas) o Pub/Sub. Preserva los límites de las transacciones y los metadatos de confirmación (commit); no se garantiza el orden global. Aplique la ordenación downstream por clave y marca de tiempo de confirmación. Espere una entrega at-least-once (al menos una vez); diseñe consumidores idempotentes (p. ej., MERGE de BigQuery con identificadores de cambio).
- Database Migration Service (DMS): Para migraciones de bases de datos con un tiempo de inactividad mínimo usando replicación nativa. DMS crea una instantánea (snapshot) consistente y luego replica continuamente los cambios usando GTID/LSN/SCN. Está diseñado específicamente para migraciones lift-and-shift, no para transformaciones arbitrarias. Para análisis, complemente DMS con Dataflow o Data Fusion si es necesario.
- Cloud Data Fusion: Un servicio de integración gestionado con conectores a sistemas relacionales, SaaS, de archivos y de mensajería. Construya canalizaciones (pipelines) con etapas de transformación (uniones, agregaciones, conversiones de formato, recetas personalizadas de Wrangler) y capture el linaje a través de fuentes y campos. Operacionalmente, programa, reintenta y emite métricas. Use Data Fusion para ELT/ETL sin código o con poco código (no/low-code) y para centralizar la gestión de conectores.
- Storage Transfer Service (STS): Transferencias gestionadas y programadas desde AWS S3, Azure Blob, on-prem (usando agentes), SFTP y listas de URL a Cloud Storage. Soporta manifiestos, sincronización incremental, control del ancho de banda e integridad verificada por suma de comprobación (checksum). Los modos de fallo incluyen la ineficiencia con archivos pequeños y la limitación de peticiones de la API (API throttling); mitíguelo con el procesamiento por lotes y una concurrencia ajustable.
- Transfer Appliance: Dispositivo (appliance) offline y cifrado para la siembra inicial de datos a escala de múltiples terabytes a petabytes cuando el ancho de banda de la red es limitado o los datos son demasiado sensibles para un tránsito prolongado. La cadena de custodia y el cifrado están incorporados. Después de la siembra, continúe con STS o CDC para los deltas.
- Cloud Pub/Sub + Dataflow: Pub/Sub desacopla productores y consumidores para patrones de streaming o de microlotes. Dataflow ofrece procesamiento de flujos (stream)/lotes con estado, autoescalado, con puntos de control (checkpointing) y marcas de agua (watermarking). Use la BigQuery Storage Write API para streaming de baja latencia con garantías exactly-once (exactamente una vez) por flujo predeterminado; de lo contrario, confíe en la semántica de deduplicación de
insertId.
Para migraciones de Hadoop a Dataproc, minimice el uso de Persistent Disk almacenando los datos en Cloud Storage con el conector de GCS y use clústeres efímeros o con autoescalado. Esto evita los grandes costes de almacenamiento en bloque (block storage) al tiempo que preserva la semántica compatible con HDFS para el procesamiento.
Esquema, validación y calidad de los datos en el límite
- Mapeo de esquemas y conversión de tipos: estandariza a esquemas fuertemente tipados desde el principio. Avro o Parquet preservan el esquema y evolucionan de forma limpia. En BigQuery, prefiere tablas particionadas y en clústeres para reducir el costo de escaneo. Ejemplo: crear una tabla particionada para análisis diario CREATE TABLE dataset.tracking_table ( event_ts TIMESTAMP, device_id STRING, payload STRING ) PARTITION BY DATE(event_ts) CLUSTER BY device_id;
- Manejo de registros malformados: dirige los rechazos a una cola de mensajes no entregados (dead-letter queue) (Pub/Sub) o a un bucket de cuarentena en Cloud Storage. Usa salidas secundarias (side outputs) en Dataflow o recolectores de errores en Data Fusion. Registra los errores de análisis con cargas útiles de muestra y versiones de esquema para su clasificación.
- Validación: realiza comprobaciones de límites antes de la persistencia:
- Estructural: conformidad del esquema, campos obligatorios, tipos de datos, dominios de enumeración.
- Referencial: existencia de claves foráneas mediante búsquedas en dimensiones cacheadas.
- Razonabilidad: rangos para marcas de tiempo, geocercas, montos no negativos.
- Unicidad: colisiones de clave primaria o clave compuesta.
- Carga idempotente: utiliza claves determinísticas y operaciones de inserción/actualización (upsert). En BigQuery, implementa MERGE con una clave de cambio natural o subrogada. Ejemplo: MERGE dataset.orders T USING dataset.orders_stage S ON T.order_id = S.order_id WHEN MATCHED THEN UPDATE SET amount = S.amount, status = S.status WHEN NOT MATCHED THEN INSERT (order_id, amount, status) VALUES (S.order_id, S.amount, S.status);
- Marcas de agua y retraso (Watermarking and lateness): en los pipelines de streaming, configura las marcas de agua de tiempo de evento (event-time watermarks) y el retraso permitido (allowed lateness) para equilibrar la completitud y la latencia. Los datos tardíos se dirigen a rutas correctivas o activan rellenos de datos (backfills).
- Reconciliación: rastrea los recuentos de filas y las sumas de verificación (checksums) por partición/ventana desde el origen hasta el destino. Captura las posiciones del registro de CDC (LSN/SCN) y las marcas de tiempo de confirmación (commit); almacénalas en una tabla de control para probar la continuidad e identificar brechas.
Conectividad, confiabilidad y operaciones
Conectividad de red y acceso privado:
- Híbrido: usa Cloud VPN o Dedicated/Partner Interconnect para conectividad privada. Habilita Private Google Access o Private Service Connect para el acceso privado a las API de Google como Cloud Storage.
- Seguridad: usa cuentas de servicio para la identidad de las cargas de trabajo, IAM con privilegios mínimos, VPC Service Controls para la prevención de la exfiltración de datos y CMEK donde sea necesario.
- Rendimiento (Throughput): escala el paralelismo en el cliente, pero en última instancia el ancho de banda gobierna el rendimiento. Para transferencias masivas, prefiere Transfer Appliance para la carga masiva inicial, y luego STS o CDC para las actualizaciones incrementales.
Puntos de control y contrapresión (Checkpoints and backpressure): Dataflow gestiona los puntos de control y el autoescalado; diseña destinos (sinks) que puedan absorber ráfagas (almacenando en búfer en Cloud Storage, escrituras por lotes en BigQuery). Para Pub/Sub, ajusta el control de flujo y los plazos de confirmación (ack deadlines) para evitar tormentas de reentrega de mensajes.
Ordenamiento y consistencia con CDC:
- Datastream preserva el orden dentro de la transacción y emite metadatos de confirmación (commit); los consumidores reconstruyen el orden por clave usando las marcas de tiempo de confirmación. Espera una entrega de “al menos una vez” (at-least-once); construye idempotencia.
- DMS garantiza la consistencia de la base de datos durante la transición de la instantánea (snapshot) y la replicación utilizando registros nativos. Usa réplicas de lectura o estrategias de escritura dual para una transición por fases.
Estrategia de archivos para analítica: para acceso de gran volumen y con múltiples motores, almacena los datos canónicos en Cloud Storage y, donde sea rentable, expón tablas externas permanentes para consultas ad hoc. Para la analítica de producción, carga los datos en tablas particionadas de BigQuery para minimizar el costo de escaneo por consulta.
Optimización de archivos pequeños: agrupa los archivos pequeños (p. ej., ~1,000 por archivo tar) antes de la transferencia y luego expándelos en la nube. Usa gsutil en paralelo y reglas de ciclo de vida para organizar por niveles y expirar los artefactos de preparación (staging).
Riesgos operativos y mitigaciones:
- Deriva de esquema desde SaaS: habilita la evolución del esquema en Data Fusion y exige compatibilidad. Alerta sobre cambios que rompan la compatibilidad.
- Zona horaria y codificación: normaliza a UTC y UTF-8 en la ingesta.
- Brechas en CDC: monitorea la retención de registros del origen; alerta cuando el retraso de la réplica se acerca a los límites de retención.
- Cuotas: inserción por streaming de BigQuery, límites de tasa de API; procesa por lotes cuando te acerques a los límites.
Transición (Cutover), Carga Histórica (Backfill) y Verificación
- Planificación de la transición (cutover):
- Big bang: congelación corta, cambio único. Menor complejidad operativa; mayor riesgo si se necesita una reversión (rollback).
- Por fases o azul/verde (blue/green): ejecución dual con escrituras duplicadas (mirrored writes), desvío progresivo de tráfico y lecturas en segundo plano (shadow reads). Mayor costo; reversión (rollback) más segura.
- Carga histórica (Backfill):
- Realizar una carga masiva inicial (Transfer Appliance o STS) usando Avro/Parquet para preservar el esquema. Particionar y agrupar en clústeres durante la carga para evitar reprocesos.
- Iniciar CDC en una posición de log conocida concurrente con la instantánea (snapshot) para capturar los deltas durante la transferencia masiva. Reconciliar en una marca de agua (watermark) común antes de abrir a producción.
- Reversión (Rollback):
- Mantener el sistema heredado (legacy) en modo de solo lectura durante la verificación. Para escenarios de escritura dual, controlar las escrituras detrás de una bandera de funcionalidad (feature flag) para revertir rápidamente. Mantener un punto de control (checkpoint) consistente para reproducir o deshacer los cambios de CDC si es necesario.
- Verificación de la migración:
- Estructural: los recuentos de filas y las sumas de verificación (checksums) por partición coinciden; el esquema y las restricciones son equivalentes.
- Temporal: sin brechas desde el límite de la instantánea (snapshot) hasta la transición (cutover); las posiciones de CDC son continuas.
- Paridad de negocio: comparar agregados y KPIs en ventanas de tiempo; ejecutar consultas de aceptación.
- Rendimiento: validar el rendimiento de la ingesta, la latencia de las consultas y el costo frente a los presupuestos.
Escenario de un Problema Práctico
Northstar Retail debe consolidar una mezcla global de sistemas transaccionales on-premise de Oracle y MySQL, eventos de un CRM SaaS y entregas diarias de archivos CSV en Google Cloud para potenciar análisis y machine learning casi en tiempo real. También necesitan migrar un clúster Hadoop heredado sin incurrir en altos gastos de almacenamiento en bloque, y lograr una transición (cutover) con tiempo de inactividad nulo o mínimo.
- Establecer conectividad híbrida segura
- Usar Partner Interconnect para el ancho de banda principal y Cloud VPN como respaldo. Habilitar Private Google Access para que las cargas de trabajo on-premise puedan acceder a Cloud Storage y Pub/Sub de forma privada. Justificación: Las rutas privadas minimizan la exposición de salida (egress) y la latencia, y Private Google Access evita los requisitos de IP públicas al tiempo que cumple con la política de seguridad.
- Cargar datos históricos de manera eficiente
- Para 800 TB de datos históricos de HDFS, copiar a Cloud Storage usando Transfer Appliance (carga masiva inicial). Después de la carga inicial, ejecutar Storage Transfer Service diariamente desde la exportación NFS on-premise para recoger los cambios hasta la transición (cutover). Justificación: Transfer Appliance evita la saturación prolongada de la red; STS proporciona una sincronización incremental programada y con suma de verificación (checksum). Almacenar en Cloud Storage con el conector GCS permite el procesamiento con Dataproc sin necesidad de 50 TB de Persistent Disk por nodo.
- Migrar bases de datos operativas con CDC
- Usar DMS para migrar MySQL y PostgreSQL con un tiempo de inactividad mínimo. Para el CDC de Oracle a sistemas de análisis, usar Datastream para depositar los datos en Cloud Storage (landing), y luego una plantilla de Dataflow proporcionada por Google para cargar en BigQuery. Justificación: DMS aprovecha la replicación nativa para una sincronización continua y de instantáneas (snapshot) fiable; Datastream proporciona CDC sin servidor (serverless) con metadatos de confirmación (commit), mientras que la plantilla de Dataflow asegura escrituras ordenadas e idempotentes en BigQuery.
- Ingerir fuentes de datos SaaS y basadas en archivos
- Construir pipelines de Cloud Data Fusion usando conectores SaaS para eventos del CRM con tokens incrementales, y un pipeline de archivos para ingerir los CSV diarios desde un SFTP de un proveedor a través de STS. Normalizar a formato Avro en un bucket de Cloud Storage curado, y luego cargar en tablas particionadas de BigQuery. Justificación: Data Fusion centraliza los conectores, la transformación y el linaje de datos. Estandarizar en Avro preserva el esquema y facilita su evolución; las tablas particionadas de BigQuery reducen el costo de las consultas.
- Procesar eventos en tiempo real (streaming)
- Publicar eventos web y de tiendas en Pub/Sub. Procesar con Dataflow para el análisis sintáctico (parsing), validación, enriquecimiento y asignación de marcas de agua (watermarking); escribir en BigQuery a través de la Storage Write API y archivar los datos brutos en formato Avro en Cloud Storage. Justificación: Pub/Sub desacopla productores y consumidores; Dataflow proporciona autoescalado, procesamiento con estado (stateful), puntos de control (checkpoints) y manejo de datos tardíos; la escritura dual asegura tanto análisis de baja latencia como retención duradera de los datos brutos.
- Aplicar controles de calidad de datos y de esquema en la frontera
- Implementar un registro de esquemas y validación en Dataflow/Data Fusion. Enrutar los registros malformados a un bucket de cuarentena en GCS y a un tema de mensajes no entregados (dead-letter topic) de Pub/Sub. Aplicar verificaciones de dominio (p. ej., códigos de moneda, marcas de tiempo UTC) y deduplicar usando claves compuestas. Justificación: El rechazo temprano y la cuarentena evitan que los datos erróneos se propaguen; la idempotencia y la deduplicación protegen contra la entrega “al menos una vez” (at-least-once) de las fuentes de CDC y streaming.
- Optimizar el almacenamiento y el acceso para análisis
- Cargar los conjuntos de datos curados en tablas de BigQuery particionadas y agrupadas en clústeres. Exponer los archivos brutos como tablas externas permanentes para exploración de baja frecuencia. Para las cargas de trabajo OLTP que siguen siendo transaccionales, mantener Cloud SQL con réplicas de lectura. Justificación: El particionamiento y la agrupación en clústeres minimizan el costo de escaneo; las tablas externas evitan cargas innecesarias para accesos ocasionales; Cloud SQL preserva la semántica ACID para las aplicaciones transaccionales.
- Planificar la transición (cutover), la carga histórica (backfill) y la reversión (rollback)
- Ejecutar instantánea (snapshot) + CDC para cada RDBMS; alcanzar un punto de reconciliación donde los recuentos de filas y las sumas de verificación (checksums) coincidan. Ejecutar en modo azul/verde (blue/green) con escrituras duales durante 48 horas, desviando las lecturas a BigQuery gradualmente. Mantener una bandera de funcionalidad (feature flag) para revertir las escrituras si se detectan discrepancias. Justificación: El modo azul/verde (blue/green) reduce el riesgo; la verificación en una marca de agua (watermark) conocida asegura la completitud; las banderas de funcionalidad permiten una reversión (rollback) rápida.
- Verificación y observabilidad
- Construir tablas de control que capturen el LSN/SCN de origen, las marcas de tiempo de confirmación (commit), los recuentos de filas y las sumas de verificación (checksums) por partición. Monitorear el retraso (lag) de Datastream, el estado de replicación de DMS, las marcas de agua (watermarks) de Dataflow, la acumulación de mensajes (backlog) de Pub/Sub, el estado de los trabajos de STS y las métricas de inserción por streaming de BigQuery. Justificación: El linaje de datos de extremo a extremo y los controles cuantitativos proporcionan una prueba auditable de la corrección y alertas oportunas sobre brechas o retrasos (lag).
Al separar las capas de aterrizaje (landing), curación y servicio (serving); usar Cloud Storage como un área de preparación (staging) y archivo duradero y de bajo costo; aprovechar DMS/Datastream para CDC con consumidores idempotentes; y aplicar controles de esquema y calidad en la entrada (ingress), Northstar Retail logra una ingesta segura y escalable y una migración verificable de bajo riesgo con un costo predecible.
← Spark · Todos los dominios · Orquestación de Flujos de Trabajo y Automatización de Canalizaciones →
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 →