Google PDE: Spark, Dataproc y Procesamiento de Datos Distribuido — 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.
Información general
Apache Spark en Google Cloud Dataproc proporciona una plataforma elástica y administrada para el procesamiento distribuido de datos. Puedes elegir entre clústeres de Dataproc de larga duración o efímeros y Dataproc Serverless para Spark, dependiendo de las necesidades de control, la variabilidad del tiempo de ejecución y la sobrecarga de gestión. Spark ofrece abstracciones resilientes (RDDs), APIs relacionales (DataFrames y Spark SQL) y un motor de ejecución de DAG tolerante a fallos optimizado para ETL iterativo y por lotes a escala. En Google Cloud, Cloud Storage reemplaza a HDFS para un almacenamiento duradero y de bajo costo; el conector de BigQuery permite la descarga analítica directa; y Dataproc Metastore centraliza la gestión de esquemas. Las soluciones eficaces alinean los ciclos de vida del almacenamiento y el cómputo, ajustan Spark a la carga de trabajo, instrumentan la observabilidad y aplican seguridad con el mínimo privilegio y aislamiento de red.
Arquitectura de Dataproc: Clústeres, Serverless, Almacenamiento y Metastore
- Tipos de clúster y roles de nodo
- Los nodos primarios (maestros) alojan YARN, el NameNode de HDFS (si se usa) y las interfaces de usuario del driver de Spark; el modo de alta disponibilidad (HA) utiliza múltiples primarios.
- Los nodos de trabajo (workers) ejecutan los ejecutores y los DataNodes de HDFS (si se usan).
- Los workers secundarios/auxiliares suelen ser preemptibles/spot para obtener capacidad elástica de menor costo sin roles de HDFS.
- Las imágenes empaquetan versiones del sistema operativo y de los componentes (por ejemplo, 2.1-debian11, 2.2-ubuntu20); fija las versiones de las imágenes para controlar la compatibilidad de Spark/Hadoop y actualizar de forma deliberada.
- Component Gateway publica las interfaces de usuario (Spark History Server, YARN RM) de forma segura a través de HTTPS.
- Dataproc Serverless para Spark
- Sin aprovisionamiento de clústeres, con autoescalado automático y facturación por segundo para ejecutores y drivers. Ideal para trabajos esporádicos o con picos de carga, o cuando se busca minimizar la sobrecarga operativa.
- Desventajas: menos controles de bajo nivel que los clústeres; la latencia de inicio del trabajo puede ser mayor que en clústeres ya activos (warm); utiliza métricas y registros de eventos serverless para la solución de problemas.
- Autoescalado
- Las políticas de autoescalado de clústeres añaden/eliminan workers basándose en métricas de YARN/Spark y periodos de enfriamiento (cooldowns), ajustando por separado los grupos de workers primarios y secundarios.
- El autoescalado serverless es gestionado por el servicio; diseña para que sea paralelo por partición y evita cuellos de botella serializados para obtener el mejor escalado.
- Almacenamiento y conectores
- Prefiere Google Cloud Storage (GCS) como el sistema de registro (system-of-record); desacopla el cómputo del almacenamiento, reduce el costo de los discos persistentes y sobrevive a los ciclos de vida del clúster.
- El conector de GCS (gs://) se integra con Hadoop/Spark. Las escrituras en almacenes de objetos utilizan protocolos de confirmación (commit); establece el algoritmo v2 de FileOutputCommitter para reducir la sobrecarga de renombrado y acelerar las confirmaciones de trabajos en GCS:
--conf mapreduce.fileoutputcommitter.algorithm.version=2
```
- Usa Parquet/ORC con poda de columnas (*column pruning*) y empuje de predicados (*predicate pushdown*). Gestiona los archivos pequeños mediante compactación para apuntar a un tamaño de 128–512 MiB por archivo para un escaneo eficiente.
- Metastore de Hive
- Centraliza los esquemas y metadatos de las tablas en Dataproc Metastore (un Apache Hive Metastore administrado) o en un metastore respaldado por Cloud SQL para compartir catálogos entre clústeres.
- Utiliza tablas externas que apunten a GCS para mayor durabilidad; particiona por fecha/hora para acotar el costo del escaneo.
- Trabajos, inicialización y flujos de trabajo
- Envía trabajos de spark, pyspark, spark-sql o hadoop. Las acciones de inicialización instalan bibliotecas o agentes adicionales en la creación del clúster (por ejemplo, conectores, bibliotecas de Python).
- Las plantillas de flujo de trabajo (*Workflow templates*) parametrizan *pipelines* de varios pasos; pueden crear clústeres efímeros por flujo de trabajo y luego eliminarlos. Esto mejora el aislamiento y reduce el costo por inactividad.
- Se recomiendan clústeres efímeros para ETL por lotes; los datos y el metastore residen fuera del clúster (en GCS, Dataproc Metastore, BigQuery).
- Integración con BigQuery
- El conector de Spark para BigQuery lee/escribe directamente en BigQuery; considera la BigQuery Storage Read API para mayor rendimiento (*throughput*) y la Write API para inserciones de *streaming* de menor latencia y con semántica *exactly-once*.
- Para el mantenimiento de tablas, realiza operaciones MERGE o sobrescrituras de particiones posteriores (*downstream*) en BigQuery para finalizar las cargas de forma atómica.
### Modelo de Spark, ajuste de rendimiento y fiabilidad
- APIs y ejecución
- RDDs: de bajo nivel, inmutables, con seguridad de tipos en Scala/Java; tú controlas la partición y la persistencia.
- DataFrames/Datasets: relacionales, optimizados por Catalyst; prefiérelos para ETL debido a la optimización de consultas y la generación de código.
- Las transformaciones son perezosas (*lazy*) (map, filter, join); las acciones desencadenan la ejecución (count, collect, save). Spark construye un DAG de etapas (*stages*) divididas por *shuffles*; las tareas se ejecutan por partición.
- Partición y *shuffle*
- Partición de entrada: suficientes particiones para utilizar todos los núcleos; comienza con 2 a 4 veces el total de núcleos de los ejecutores (*executors*). Contrólalo a través de spark.default.parallelism (para RDDs) y las opciones del lector (para DataFrames).
- Particiones de *shuffle*: el valor predeterminado de 200 a menudo subaprovisiona o sobreaprovisiona. Ajústalo:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
- Apunta a ~100–256 MiB por partición después de transformaciones amplias (wide transforms); un tamaño demasiado pequeño causa sobrecarga en el planificador (scheduler) y uno demasiado grande arriesga un OOM en el ejecutor.
- El shuffle es el costo dominante para joins, groupBy y orderBy. Asegura una memoria y disco adecuados para el ejecutor; considera SSDs locales para shuffles pesados en clústeres.
- Asimetría (skew) y estrategia de join
- Detecta la asimetría (skew) (tiempos de ejecución de tareas de cola larga, tamaños de partición grandes). Mitigaciones:
- Haz broadcast de tablas pequeñas para evitar shuffles:
- Detecta la asimetría (skew) (tiempos de ejecución de tareas de cola larga, tamaños de partición grandes). Mitigaciones:
--conf spark.sql.autoBroadcastJoinThreshold=64m
```
- Añade *salt* a las claves para particiones calientes (*hot partitions*); aplica preagregación en el lado del *map*; filtra temprano.
- Habilita Adaptive Query Execution (AQE) para fusionar particiones post-*shuffle* y manejar *joins* asimétricos (*skewed joins*):
--conf spark.sql.adaptive.enabled=true
```
- Caché, checkpointing y linaje
- Almacena en caché los DataFrames intermedios calientes (hot) con moderación cuando se reutilizan; prefiere MEMORY_AND_DISK para evitar OOM.
- Haz checkpoint de linajes largos en GCS o HDFS para acotar el recálculo en caso de fallos.
- Ejecutores (executors) y asignación dinámica
- Dimensiona correctamente los ejecutores para equilibrar el paralelismo y la sobrecarga del GC:
- Núcleos por ejecutor: 2–5 para tareas equilibradas de E/S y CPU; menos núcleos reducen las pausas del GC.
- Sobrecarga de memoria: establece spark.yarn.executor.memoryOverhead para shuffles amplios.
- Habilita la asignación dinámica con el servicio de shuffle externo en clústeres para escalar los ejecutores según la carga de trabajo:
- Dimensiona correctamente los ejecutores para equilibrar el paralelismo y la sobrecarga del GC:
--conf spark.dynamicAllocation.enabled=true
--conf spark.shuffle.service.enabled=true
--conf spark.dynamicAllocation.minExecutors=0
--conf spark.dynamicAllocation.maxExecutors=200
```
- Patrones de tolerancia a fallos para ETL por lotes (*batch*)
- Escrituras idempotentes: escribe en una ruta temporal o de *staging*, luego promueve atómicamente con un *commit* a nivel de directorio; para BigQuery, escribe en una tabla de *staging* y usa MERGE:
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```
- Procesamiento incremental: usa filtrado basado en marcas de agua (watermarks) en las particiones de ingestion_date; mantén un manifiesto de procesados en GCS para evitar el reprocesamiento.
- Manejo de mensajes fallidos (dead-letter): en errores de análisis o validación, deriva los registros erróneos a una ruta o tabla de cuarentena con diagnósticos. Para una aplicación estricta del esquema y DLQs integradas, considera Dataflow; con Spark, implementa un try/catch por registro y un sink separado.
Seguridad, observabilidad y costo
- Identidad y acceso
- Ejecute clústeres y trabajos con cuentas de servicio dedicadas y con el principio de privilegio mínimo (least-privilege) en IAM. Otorgue solo los roles necesarios, por ejemplo:
- roles/dataproc.worker a las cuentas de servicio de las instancias
- roles/storage.objectViewer u objectAdmin para las rutas de E/S de GCS
- roles/bigquery.dataEditor en los datasets de destino
- Para Dataproc Serverless, utilice cuentas de servicio por trabajo para delimitar el acceso.
- Ejecute clústeres y trabajos con cuentas de servicio dedicadas y con el principio de privilegio mínimo (least-privilege) en IAM. Otorgue solo los roles necesarios, por ejemplo:
- Aislamiento de red y cifrado
- Use clústeres con IP privada en una subred de VPC, restrinja con firewall las IU del nodo maestro y habilite el Acceso Privado a Google para GCS/BigQuery sin salida pública.
- Coloque los clústeres en proyectos de VPC Compartida (Shared VPC) para un control centralizado. Opcionalmente, habilite Kerberos en Dataproc para la autenticación dentro del clúster.
- Cifre en reposo con CMEK: configure CMEK en los buckets de GCS, Persistent Disks, Dataproc Metastore y BigQuery; use TLS en tránsito por defecto.
- Registros, historial y métricas
- Habilite los registros de eventos de Spark en GCS y despliegue el History Server:
--conf spark.eventLog.enabled=true
--conf spark.eventLog.dir=gs://bucket/spark-events/
```
- Dataproc transmite los registros del driver y de YARN a Cloud Logging; expórtelos a receptores (sinks) para retención/análisis forense.
- Supervise con las métricas de Cloud Monitoring: contenedores pendientes de YARN, CPU, memoria, estado de HDFS (si se usa), rendimiento de GCS. Alerte sobre reintentos de etapa prolongados, pérdida de ejecutores y picos de ejecución especulativa.
- Análisis de fallos: las causas comunes incluyen nodos rezagados (stragglers) inducidos por sesgo de datos (skew), OOM (Out of Memory) del ejecutor durante el shuffle, fallos de confirmación (commit) en el almacén de objetos y pérdida de nodos interrumpibles/spot. Aumente el número de reintentos con prudencia; un exceso de reintentos puede amplificar el costo y el retraso.
- Optimización de costos
- Use clústeres efímeros o Dataproc Serverless para evitar el costo por inactividad; mantenga los datos en GCS para minimizar el uso de discos persistentes.
- Añada trabajadores secundarios interrumpibles/spot para absorber la demanda máxima; diseñe para el recálculo, ya que las tareas en nodos perdidos se reintentan. No coloque nodos maestros en nodos interrumpibles.
- Dimensione correctamente los tipos de máquina y use el autoescalado para reducir la capacidad cuando las colas estén vacías. Prefiera Parquet/ORC con poda de particiones (partition pruning) para reducir el costo de escaneo y el uso de CPU.
- Evite los archivos pequeños compactando las salidas; menos archivos y más grandes reducen la sobrecarga de metadatos y el tiempo de ejecución del trabajo.
- Para trabajos cortos y periódicos (por ejemplo, un ETL de Spark de 30 minutos semanalmente), los trabajadores interrumpibles o el modo serverless suelen ofrecer el mejor perfil de costo.
#### Escenario de un problema práctico
Acme Retail está migrando un clúster de Hadoop on-premise de 30 nodos que ejecuta ETL nocturnos de Spark y Hive que alimentan los sistemas de análisis posteriores (downstream). Quieren reutilizar los trabajos existentes con cambios mínimos, evitar la gestión de clústeres a tiempo completo, persistir los datos más allá del ciclo de vida del clúster y reducir el costo de almacenamiento.
Enfoque:
1) Depositar datos y metadatos en servicios gestionados
- Almacene todos los datos brutos y curados en Cloud Storage usando Parquet con particionamiento (por ejemplo, dt=YYYY-MM-DD).
- Justificación: GCS es duradero, de bajo costo y desacopla el cómputo del almacenamiento, por lo que los clústeres efímeros y los trabajos serverless pueden ejecutarse sin discos persistentes. El formato Parquet particionado permite el empuje de predicados (predicate pushdown) y escaneos eficientes.
2) Centralizar el catálogo con Dataproc Metastore
- Migre el metastore de Hive a Dataproc Metastore. Cree tablas externas de Hive que hagan referencia a rutas de GCS y conserve la lógica de esquema/partición existente.
- Justificación: Un metastore gestionado permite que múltiples clústeres efímeros y trabajos serverless compartan definiciones de tablas sin necesidad de ejecutar una instancia de MySQL/PostgreSQL de alta disponibilidad (HA).
3) Usar clústeres efímeros de Dataproc para ETL por lotes y plantillas de flujo de trabajo para la orquestación
- Defina una plantilla de flujo de trabajo que cree un clúster con la imagen requerida (por ejemplo, 2.1-debian11), ejecute trabajos de Spark (spark-sql y pyspark) y elimine el clúster al finalizar. Añada acciones de inicialización para instalar cualquier biblioteca personalizada.
- Justificación: Los clústeres efímeros eliminan el costo por inactividad y aíslan las dependencias de los trabajos. Las plantillas de flujo de trabajo proporcionan repetibilidad y parametrización (fechas, rutas de entrada).
4) Habilitar el autoescalado y los trabajadores interrumpibles
- Asocie una política de autoescalado con un pequeño grupo de trabajadores principales (core) y un grupo más grande de trabajadores secundarios interrumpibles; ajuste los períodos de enfriamiento (cooldowns) para reducir la escala rápidamente después de la ejecución.
- Justificación: Los trabajadores principales mantienen la estabilidad del clúster; los trabajadores interrumpibles absorben los shuffles y las transformaciones amplias (wide transformations) a un costo menor. Los reintentos de Spark/YARN gestionan las tareas perdidas por interrupción.
5) Integrar con BigQuery a través del conector de Spark para BigQuery
- Para las cargas de dimensiones/hechos, escriba los resultados de Spark en tablas de staging en BigQuery y luego ejecute sentencias MERGE para actualizar los destinos de forma atómica. Cuando la sobrescritura directa sea segura, escriba en tablas particionadas usando el modo de sobrescritura de partición.
- Justificación: BigQuery sirve para análisis y BI a escala; la combinación de staging+MERGE produce operaciones de tipo upsert transaccionales desde lotes de Spark, reduciendo la inconsistencia en los sistemas posteriores.
6) Ajustar Spark para rendimiento y fiabilidad
- Establezca las particiones de shuffle en relación con los núcleos de los ejecutores y habilite AQE:
--conf spark.sql.shuffle.partitions=600
--conf spark.sql.adaptive.enabled=true
```
- Use broadcast joins para dimensiones pequeñas y haga checkpoint de linajes largos en GCS para mayor estabilidad.
- Justificación: Un particionamiento adecuado reduce el sesgo de datos (skew) y la sobrecarga del planificador (scheduler); AQE se adapta a los perfiles de datos en tiempo de ejecución; el checkpointing limita el recálculo después de fallos.
Reforzar la seguridad y las redes
- Ejecute clústeres con cuentas de servicio dedicadas que otorguen solo los roles necesarios para las rutas de GCS, el metastore y los datasets de BigQuery. Cree clústeres con IP privada en una subred restringida con Acceso Privado a Google y limite el acceso a la IU mediante reglas de firewall.
- Justificación: El principio de privilegio mínimo y el aislamiento de red reducen la superficie de ataque; el egreso del plano de control privado evita la exposición pública.
Instrumentar registros, historial y alertas
- Habilite los registros de eventos de Spark en GCS y despliegue el History Server; dirija los registros del driver/YARN a Cloud Logging con retención. Añada alertas de Monitoring para contenedores pendientes durante mucho tiempo, fallos de tareas repetidos o duración excesiva del trabajo.
- Justificación: Los registros centralizados apoyan el análisis de causa raíz; las alertas proactivas detectan sesgos, OOMs o E/S degradada de forma temprana.
Modernizar selectivamente con Dataproc Serverless para picos ad hoc y elásticos
- Mueva las cargas de trabajo esporádicas o exploratorias de Spark SQL a Dataproc Serverless; mantenga las canalizaciones nocturnas en clústeres efímeros hasta que estén completamente validadas en serverless.
- Justificación: Serverless elimina las operaciones del clúster y escala automáticamente, ideal para cargas impredecibles; los flujos de trabajo existentes continúan con cambios mínimos en el código.
Validar los committers del almacén de objetos y la gestión de archivos pequeños
- Establezca el algoritmo v2 de FileOutputCommitter y compacte las salidas a 256–512 MiB por archivo mediante repartition/coalesce antes de las escrituras.
- Justificación: Los almacenes de objetos carecen de renombrado atómico; los committers optimizados reducen la sobrecarga de copia/renombrado. La compactación mitiga el problema de los archivos pequeños para mejorar el rendimiento y el costo.
Este diseño reutiliza los trabajos existentes de Spark y Hive con una refactorización mínima, asegura la durabilidad de los datos en GCS, centraliza los esquemas, contiene el radio de impacto de seguridad, proporciona una observabilidad robusta y optimiza el costo a través de clústeres efímeros, autoescalado, capacidad interrumpible y el uso específico de la ejecución serverless.
← Mensajería · Todos los dominios · Ingesta →
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 →