Google PDE: Spark、Dataproc 与分布式数据处理 — 学习指南
属于 Google Professional Data Engineer — 学习指南. 使用经过验证的答案练习: Google 考试中心, 或参加限时模拟考试: ExamRoll.io.
概览
Google Cloud Dataproc 上的 Apache Spark 为分布式数据处理提供了一个托管的弹性平台。您可以根据控制需求、运行时可变性和管理开销,在长期运行的或临时性的 Dataproc 集群与 Dataproc Serverless for Spark 之间进行选择。Spark 提供了弹性抽象 (RDDs)、关系型 API (DataFrames 和 Spark SQL) 以及一个容错的 DAG 执行引擎,该引擎为大规模的迭代式和批量 ETL 进行了优化。在 Google Cloud 上,Cloud Storage 取代 HDFS 作为持久、低成本的存储;BigQuery 连接器支持直接进行分析负载卸载;Dataproc Metastore 则集中化了模式管理。有效的解决方案会使存储和计算的生命周期保持一致,根据工作负载调优 Spark,实施可观测性,并应用最小权限和网络隔离的安全措施。
Dataproc 架构:集群、无服务器、存储和元存储
- 集群类型和节点角色
- 主(master)节点托管 YARN、HDFS NameNode(如果使用)和 Spark 驱动程序 UI;HA(高可用性)模式使用多个主节点。
- 工作器节点运行执行器和 HDFS DataNode(如果使用)。
- 辅助/次要工作器通常是可抢占/Spot 实例,用于提供弹性的、成本更低的容量,不承担 HDFS 角色。
- 镜像捆绑了操作系统和组件版本(例如,2.1-debian11、2.2-ubuntu20);固定镜像版本以控制 Spark/Hadoop 兼容性并进行审慎升级。
- 组件网关通过 HTTPS 安全地发布 UI(Spark History Server、YARN RM)。
- Dataproc Serverless for Spark
- 无需集群预配、自动进行自动扩缩容,并对执行器和驱动程序按秒计费。非常适合零星或突发作业,或需要最大限度减少运维开销的场景。
- 权衡:与集群相比,可供调整的底层选项较少;作业启动延迟可能高于已预热的集群;使用无服务器指标和事件日志进行故障排查。
- 自动扩缩容
- 集群自动扩缩容策略根据 YARN/Spark 指标和冷却时间来增减工作器,并分别对主工作器组和辅助工作器组进行调优。
- 无服务器的自动扩缩容由服务托管;设计时应使其分区并行,并避免序列化瓶颈,以实现最佳扩缩容效果。
- 存储和连接器
- 首选 Google Cloud Storage (GCS) 作为记录系统;它将计算与存储解耦,降低了永久性磁盘成本,并且能在集群生命周期结束后继续存在。
- GCS 连接器 (gs://) 与 Hadoop/Spark 集成。向对象存储的写入操作使用提交协议;将 FileOutputCommitter 算法设置为 v2,以减少重命名开销并加速在 GCS 上的作业提交:
--conf mapreduce.fileoutputcommitter.algorithm.version=2
```
- 使用 Parquet/ORC 格式,并结合列裁剪和谓词下推。通过合并(compaction)来管理小文件,目标是使每个文件大小在 128–512 MiB 之间,以实现高效扫描。
- Hive 元存储
- 将模式和表元数据集中到 Dataproc Metastore(托管的 Apache Hive Metastore)或由 Cloud SQL 支持的元存储中,以便在集群间共享目录。
- 使用指向 GCS 的外部表以确保持久性;按日期/小时进行分区以限制扫描成本。
- 作业、初始化和工作流
- 提交 `spark`、`pyspark`、`spark-sql` 或 `hadoop` 作业。初始化操作在集群创建时安装额外的库或代理(例如,连接器、Python 库)。
- 工作流模板可将多步骤流水线参数化;它们可以为每个工作流创建临时集群,然后在完成后将其拆除。这可以改善隔离性并削减空闲成本。
- 对于批量 ETL,推荐使用临时集群;数据和元存储存在于集群之外(GCS、Dataproc Metastore、BigQuery)。
- BigQuery 集成
- Spark BigQuery 连接器可直接读写 BigQuery;考虑使用 BigQuery Storage Read API 以获得高吞吐量,使用 Write API 以实现更低延迟、仅一次的流式插入。
- 对于表维护,在 BigQuery 中执行下游的 MERGE/分区覆盖操作,以原子方式完成加载。
### Spark 模型、性能调优与可靠性
- API 与执行
- RDDs:低级别、不可变、在 Scala/Java 中是类型安全的;由您控制分区和持久化。
- DataFrames/Datasets:关系型,经 Catalyst 优化;因其查询优化和代码生成能力,在 ETL 场景中应优先使用。
- 转换操作(Transformation)是惰性的(如 map、filter、join);动作操作(Action)触发执行(如 count、collect、save)。Spark 会构建一个由 shuffle 分割的阶段(stage)组成的 DAG;任务(task)在每个分区上运行。
- 分区与 shuffle
- 输入分区:确保有足够的分区来利用所有核心;初始可设置为 executor 总核心数的 2-4 倍。通过
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
(针对 RDDs)和读取器选项(针对 DataFrames)进行控制。
- Shuffle 分区:默认值 200 常常导致资源供给不足或过量。进行调优:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
- 目标是让宽转换(wide transform)后的每个分区大小在 100–256 MiB 左右;分区太小会导致调度器开销,太大则有 executor OOM 的风险。
- 对于 join、groupBy 和 orderBy 操作,Shuffle 是主要的成本开销。确保 executor 有足够的内存和磁盘空间;对于集群上的重度 shuffle 负载,可考虑使用本地 SSD。
- 数据倾斜与 join 策略
- 检测数据倾斜(任务运行时长出现长尾、分区大小差异悬殊)。缓解措施:
- 广播小表以避免 shuffle:
--conf spark.sql.autoBroadcastJoinThreshold=64m
```
- 为热点分区的数据键(key)加盐;应用 map 端预聚合;尽早进行数据过滤。
- 启用自适应查询执行(AQE)来合并 shuffle 后的分区并处理倾斜的 join:
--conf spark.sql.adaptive.enabled=true
```
- 缓存、检查点与血缘
- 当中间 DataFrame 会被重复使用时,谨慎地缓存热点数据;优先使用 MEMORY_AND_DISK 以避免 OOM。
- 将长血缘关系的数据 checkpoint 到 GCS 或 HDFS,以限制失败时需要重新计算的范围。
- Executor 与动态分配
- 合理设置 executor 的大小,以平衡并行度与 GC 开销:
- 每个 executor 的核心数:对于均衡的 I/O/CPU 任务,设置为 2-5 个;较少的核心数可以减少 GC 暂停时间。
- 内存开销:为宽 shuffle 设置 `spark.yarn.executor.memoryOverhead`。
- 在集群上启用动态分配和外部 shuffle 服务,使 executor 数量能随工作负载伸缩:
--conf spark.dynamicAllocation.enabled=true
--conf spark.shuffle.service.enabled=true
--conf spark.dynamicAllocation.minExecutors=0
--conf spark.dynamicAllocation.maxExecutors=200
```
- 批处理 ETL 的容错模式
- 幂等写入:先写入临时/暂存路径,然后通过目录级提交原子性地提升为正式数据;对于 BigQuery,则写入暂存表并执行 MERGE 操作:
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```
- 增量处理:在 `ingestion_date` 分区上使用基于水印的过滤;在 GCS 中维护一个已处理清单,以避免重复处理。
- 死信处理:当发生解析/验证错误时,将坏记录连同诊断信息分流到隔离路径/表中。若需严格的 schema 强制和内置的 DLQ(死信队列),可考虑使用 Dataflow;若使用 Spark,则需实现记录级别的 try/catch 和一个独立的 sink。
### 安全、可观测性与成本
- 身份与访问权限
- 在具有最小权限 IAM 的专用服务账号下运行集群和作业。仅授予所需角色,例如:
- roles/dataproc.worker 授予实例服务账号
- roles/storage.objectViewer 或 objectAdmin 授予 GCS I/O 路径
- roles/bigquery.dataEditor 授予目标数据集
- 对于 Dataproc Serverless,使用按作业的服务账号来限定访问范围。
- 网络隔离与加密
- 在 VPC 子网中使用私有 IP 集群,通过防火墙限制主节点 UI 的访问,并启用 Private Google Access 以便在没有公共出口的情况下访问 GCS/BigQuery。
- 将集群放置在 Shared VPC 项目中以实现集中控制。可选择在 Dataproc 上启用 Kerberos 以进行集群内身份验证。
- 使用 CMEK 进行静态加密:在 GCS 存储桶、永久性磁盘、Dataproc Metastore 和 BigQuery 上配置 CMEK;默认使用 TLS 进行传输中加密。
- 日志、历史记录与指标
- 启用 Spark 事件日志到 GCS 并部署 History Server:
--conf spark.eventLog.enabled=true
--conf spark.eventLog.dir=gs://bucket/spark-events/
```
- Dataproc 将驱动程序和 YARN 日志流式传输到 Cloud Logging;导出到接收器 (sink) 以用于保留/取证。
- 使用 Cloud Monitoring 指标进行监控:YARN 待处理容器数、CPU、内存、HDFS 健康状况(如果使用)、GCS 吞吐量。针对长时间的阶段重试、executor 丢失和推测执行峰值设置警报。
- 故障分析:常见原因包括数据倾斜导致的拖后任务 (straggler)、shuffle 期间 executor 发生 OOM (内存溢出)、对象存储提交失败以及可抢占/Spot 节点丢失。审慎地增加重试次数;过多的重试会增加成本和延迟。
- 成本优化
- 使用临时集群或 Dataproc Serverless 以避免闲置成本;将数据保存在 GCS 中以最大限度地减少永久性磁盘的使用。
- 添加可抢占/Spot 辅助工作器以吸收峰值需求;设计时应考虑可重新计算,因为丢失节点上的任务会被重试。不要将主节点放在可抢占节点上。
- 合理调整机器类型大小,并使用自动扩缩容在队列为空时缩减容量。优先使用 Parquet/ORC 格式并结合分区裁剪来降低扫描成本和 CPU 消耗。
- 通过合并输出来避免小文件;更少、更大的文件可以减少元数据开销和作业运行时间。
- 对于短期的周期性作业(例如,每周 30 分钟的 Spark ETL),可抢占工作器或无服务器模式通常能提供最佳的成本效益。
实际问题场景
Acme Retail 正在迁移一个包含 30 个节点的本地 Hadoop 集群,该集群运行夜间的 Spark 和 Hive ETL 作业,为下游分析提供数据。他们希望以最小的改动重用现有作业,避免全职管理集群,在集群生命周期结束后持久化数据,并降低存储成本。
方法:
将数据和元数据存放在托管服务中
- 使用带分区的 Parquet 格式(例如,dt=YYYY-MM-DD)将所有原始数据和整理后的数据存储在 Cloud Storage 中。
- 理由:GCS 持久、低成本,并将计算与存储解耦,因此临时集群和无服务器作业可以在没有永久性磁盘的情况下运行。分区的 Parquet 支持谓词下推和高效扫描。
使用 Dataproc Metastore 集中化目录
- 将 Hive metastore 迁移到 Dataproc Metastore。创建引用 GCS 路径的外部 Hive 表,并保留现有的模式/分区逻辑。
- 理由:托管的 metastore 允许多个临时集群和无服务器作业共享表定义,而无需运行一个高可用的 MySQL/PostgreSQL 实例。
使用临时 Dataproc 集群进行批处理 ETL,并使用工作流模板进行编排
- 定义一个工作流模板,该模板创建一个具有所需镜像(例如,2.1-debian11)的集群,运行 Spark 作业(spark-sql 和 pyspark),并在完成后删除集群。添加初始化操作以安装任何自定义库。
- 理由:临时集群消除了闲置成本并隔离了作业依赖。工作流模板提供了可重复性和参数化能力(例如日期、输入路径)。
启用自动扩缩容和可抢占工作器
- 附加一个自动扩缩容策略,配置一个小的核心工作器组和一个较大的可抢占辅助工作器池;调整冷却时间以在运行后迅速缩减。
- 理由:核心工作器维持集群稳定性;可抢占工作器以较低成本吸收 shuffle 和宽转换操作。Spark/YARN 的重试机制会处理因抢占而丢失的任务。
通过 Spark BigQuery 连接器与 BigQuery 集成
- 对于维度/事实表加载,将 Spark 结果写入 BigQuery 的暂存表,然后运行 MERGE 语句以原子方式更新目标表。在可以直接覆盖的情况下,使用分区覆盖模式写入分区表。
- 理由:BigQuery 为分析和 BI 提供大规模服务;暂存+MERGE 的方式可以实现来自批处理 Spark 的类事务性更新插入 (upsert),从而减少下游数据的不一致性。
调整 Spark 以优化性能和可靠性
- 根据 executor 内核数设置 shuffle 分区数,并启用 AQE:
--conf spark.sql.shuffle.partitions=600
--conf spark.sql.adaptive.enabled=true
```
- 对小型维度表使用广播连接,并将长血缘关系的数据检查点到 GCS 以提高稳定性。
- 理由:正确的分区可以减少数据倾斜和调度器开销;AQE 在运行时能适应数据分布;检查点机制限制了故障后需要重新计算的范围。
7) 强化安全与网络
- 使用专用服务账号运行集群,仅授予 GCS 路径、metastore 和 BigQuery 数据集所需的角色。在受限子网中创建私有 IP 集群,启用 Private Google Access,并通过防火墙规则限制对 UI 的访问。
- 理由:最小权限原则和网络隔离减少了攻击面;私有的控制平面出口避免了公网暴露。
8) 配置日志、历史记录和警报
- 启用 Spark 事件日志到 GCS 并部署 History Server;将驱动程序/YARN 日志路由到 Cloud Logging 并设置保留策略。为长时间待处理的容器、重复的任务失败或过长的作业持续时间添加 Monitoring 警报。
- 理由:集中式日志支持根本原因分析;主动警报能及早发现数据倾斜、OOM 或 I/O 性能下降等问题。
9) 通过 Dataproc Serverless 对临时和弹性峰值负载进行选择性现代化改造
- 将零星的或探索性的 Spark SQL 工作负载迁移到 Dataproc Serverless;在无服务器模式上得到充分验证之前,保持夜间管道在临时集群上运行。
- 理由:无服务器模式消除了集群运维工作并能自动扩缩容,非常适合不可预测的负载;现有工作流可以继续运行,只需极少的代码更改。
10) 验证对象存储提交器和小文件管理
- 设置 FileOutputCommitter v2 算法,并在写入前通过 repartition/coalesce 将输出文件合并为每个文件 256–512 MiB。
- 理由:对象存储缺少原子性的重命名操作;优化的提交器减少了复制/重命名开销。文件合并缓解了小文件问题,从而提升性能并降低成本。
此设计以最小的重构代价重用了现有的 Spark 和 Hive 作业,确保了数据在 GCS 中的持久性,集中了模式定义,控制了安全风险范围,提供了强大的可观测性,并通过临时集群、自动扩缩容、可抢占容量以及有针对性地使用无服务器执行来优化成本。
---
← [消息传递、事件注入与实时服务](/cn/posts/pde-messaging-ingestion/) · [所有领域](/cn/posts/google-pde-study-guide/) · [数据注入、集成与迁移](/cn/posts/pde-ingestion-migration/) →
**[练习这些题目 →](/cn/kb/google/)** · **[在 ExamRoll.io 上限时练习 →](https://www.examroll.io/?utm_source=guide&utm_medium=referral&utm_campaign=PDE)**
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.
通过考试 →