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 架构:集群、无服务器、存储和元存储

    --conf mapreduce.fileoutputcommitter.algorithm.version=2
    ```

  - 使用 Parquet/ORC 格式,并结合列裁剪和谓词下推。通过合并(compaction)来管理小文件,目标是使每个文件大小在 128512 MiB 之间,以实现高效扫描。
- Hive 元存储
  - 将模式和表元数据集中到 Dataproc Metastore(托管的 Apache Hive Metastore)或由 Cloud SQL 支持的元存储中,以便在集群间共享目录。
  - 使用指向 GCS 的外部表以确保持久性;按日期/小时进行分区以限制扫描成本。
- 作业、初始化和工作流
  - 提交 `spark``pyspark``spark-sql`  `hadoop` 作业。初始化操作在集群创建时安装额外的库或代理(例如,连接器、Python 库)。
  - 工作流模板可将多步骤流水线参数化;它们可以为每个工作流创建临时集群,然后在完成后将其拆除。这可以改善隔离性并削减空闲成本。
  - 对于批量 ETL,推荐使用临时集群;数据和元存储存在于集群之外(GCSDataproc MetastoreBigQuery)。
- BigQuery 集成
  - Spark BigQuery 连接器可直接读写 BigQuery;考虑使用 BigQuery Storage Read API 以获得高吞吐量,使用 Write API 以实现更低延迟、仅一次的流式插入。
  - 对于表维护,在 BigQuery 中执行下游的 MERGE/分区覆盖操作,以原子方式完成加载。
### Spark 模型、性能调优与可靠性

- API 与执行
  - RDDs:低级别、不可变、在 Scala/Java 中是类型安全的;由您控制分区和持久化。
  - DataFrames/Datasets:关系型,经 Catalyst 优化;因其查询优化和代码生成能力,在 ETL 场景中应优先使用。
  - 转换操作(Transformation)是惰性的(如 mapfilterjoin);动作操作(Action)触发执行(如 countcollectsave)。Spark 会构建一个由 shuffle 分割的阶段(stage)组成的 DAG;任务(task)在每个分区上运行。
- 分区与 shuffle
  - 输入分区:确保有足够的分区来利用所有核心;初始可设置为 executor 总核心数的 2-4 倍。通过
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```

(针对 RDDs)和读取器选项(针对 DataFrames)进行控制。

    --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
  ```
    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/
```

实际问题场景

Acme Retail 正在迁移一个包含 30 个节点的本地 Hadoop 集群,该集群运行夜间的 Spark 和 Hive ETL 作业,为下游分析提供数据。他们希望以最小的改动重用现有作业,避免全职管理集群,在集群生命周期结束后持久化数据,并降低存储成本。

方法:

  1. 将数据和元数据存放在托管服务中

    • 使用带分区的 Parquet 格式(例如,dt=YYYY-MM-DD)将所有原始数据和整理后的数据存储在 Cloud Storage 中。
    • 理由:GCS 持久、低成本,并将计算与存储解耦,因此临时集群和无服务器作业可以在没有永久性磁盘的情况下运行。分区的 Parquet 支持谓词下推和高效扫描。
  2. 使用 Dataproc Metastore 集中化目录

    • 将 Hive metastore 迁移到 Dataproc Metastore。创建引用 GCS 路径的外部 Hive 表,并保留现有的模式/分区逻辑。
    • 理由:托管的 metastore 允许多个临时集群和无服务器作业共享表定义,而无需运行一个高可用的 MySQL/PostgreSQL 实例。
  3. 使用临时 Dataproc 集群进行批处理 ETL,并使用工作流模板进行编排

    • 定义一个工作流模板,该模板创建一个具有所需镜像(例如,2.1-debian11)的集群,运行 Spark 作业(spark-sql 和 pyspark),并在完成后删除集群。添加初始化操作以安装任何自定义库。
    • 理由:临时集群消除了闲置成本并隔离了作业依赖。工作流模板提供了可重复性和参数化能力(例如日期、输入路径)。
  4. 启用自动扩缩容和可抢占工作器

    • 附加一个自动扩缩容策略,配置一个小的核心工作器组和一个较大的可抢占辅助工作器池;调整冷却时间以在运行后迅速缩减。
    • 理由:核心工作器维持集群稳定性;可抢占工作器以较低成本吸收 shuffle 和宽转换操作。Spark/YARN 的重试机制会处理因抢占而丢失的任务。
  5. 通过 Spark BigQuery 连接器与 BigQuery 集成

    • 对于维度/事实表加载,将 Spark 结果写入 BigQuery 的暂存表,然后运行 MERGE 语句以原子方式更新目标表。在可以直接覆盖的情况下,使用分区覆盖模式写入分区表。
    • 理由:BigQuery 为分析和 BI 提供大规模服务;暂存+MERGE 的方式可以实现来自批处理 Spark 的类事务性更新插入 (upsert),从而减少下游数据的不一致性。
  6. 调整 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.

通过考试 →

浏览 Google →

Related guides

一体化访问

一次订阅。所有考试。

所有计划均可无限制搜索答案、进行模拟测试、获取AI解释以及访问完整的资源库 — 支持20多种语言。

每月
24.87
Just €0.83/day
包含所有内容:
  • 无限答案搜索
  • 无限模拟测试
  • AI驱动的解释
  • 完整资源库
  • 20多种语言
  • 每周内容更新
  • 奖励与推荐
  • 优先支持
开始免费试用

无需信用卡*

最具价值
12个月
179.87
Just €0.49/daySave 40%
包含所有内容:
  • 无限答案搜索
  • 无限模拟测试
  • AI驱动的解释
  • 完整资源库
  • 20多种语言
  • 每周内容更新
  • 奖励与推荐
  • 优先支持
开始免费试用

无需信用卡*

✓ 包含免费计划 · ✓ 随时取消 · ✓ 所有计划均解锁完整产品