Google PDE: 数据工程架构与设计 — 学习指南
属于 Google Professional Data Engineer — 学习指南. 使用经过验证的答案练习: Google 考试中心, 或参加限时模拟考试: ExamRoll.io.
概述
Google Cloud 上的数据工程架构与设计需要在领域边界、处理模式和服务能力之间取得平衡,以交付可靠、可扩展且经济高效的数据平台。有效的设计能使存储、计算、编排和服务层可以独立扩展;将契约代码化以实现领域间的互操作;并通过可衡量的服务水平目标 (SLO) 尽早验证风险。本节总结了典型的架构风格(数据网格、数据湖、数据仓库、湖仓一体、业务型存储)、处理模式(批处理、微批处理、流处理、事件驱动、Lambda 架构),以及在可扩展性、延迟、可用性、一致性和成本之间的权衡。此外,本节还涵盖了区域和多云部署、模式演进、端到端数据生命周期、基于工作负载的服务选择,以及针对 Google Cloud 定制的风险驱动型验证实践。
架构范式与处理模式
- 数据网格、领域和数据产品:
- 授权领域团队发布“数据产品”,这些产品具有明确的所有权、SLO、访问策略和文档。使用 Dataplex 定义领域、治理元数据,并在 BigQuery 和 Cloud Storage 之间应用一致的策略。产品可以暴露 BigQuery 数据集、Pub/Sub 主题或 Cloud Storage 路径,其契约通过 Pub/Sub 模式和 BigQuery 表模式来强制执行。
- 数据湖:
- 在 Cloud Storage 中以开放格式(Parquet/Avro)存储原始数据,并配置生命周期和版本控制。适用于异构工作负载(Dataproc 上的 Spark、Dataflow、Presto/Trino)和多云可移植性。权衡:对象存储中的最终一致性语义;设计时应考虑幂等性和元数据驱动的去重机制。
- 数据仓库:
- 在 BigQuery 中进行经过整理和治理的分析。为 ANSI SQL、存储/计算分离和细粒度安全性进行了优化。权衡:流式插入在查询时会表现出短暂的数据过时;对于严格的新鲜度 SLA,应优先选择批量加载或使用带缓冲的查询进行插入。
- 湖仓一体 (Lakehouse):
- 将开放的数据湖存储与数据仓库的能力相融合。在 Google Cloud 上,将 Parquet/Avro 存储在 Cloud Storage 中;使用 BigQuery 外部表以实现经济性,使用 BigQuery 托管表以获得高性能和治理能力。Dataflow 或 Dataproc 通过分区/聚类策略来维护类似 ACID 的合并语义。
- 业务型存储架构:
- 为应用程序提供支持的低延迟事务性或键值存储。为传统 OLTP 选择 Cloud SQL,为需要水平扩展的全局一致性 SQL 选择 Cloud Spanner,为极高吞吐量的宽列访问模式选择 Bigtable。将业务型存储与分析系统分离;使用 CDC (Datastream) 将变更捕获到 Pub/Sub、Cloud Storage 或 BigQuery 中。
处理模式及其适用场景:
- 批处理:定期的、大规模的转换(例如,夜间的特征生成)。工具:Dataflow 批处理、Dataproc。失败模式:长时间运行的作业超时、数据倾斜;通过自动扩缩和重新分区来缓解。
- 微批处理:小而频繁的批次(例如,每分钟一次),以平衡新鲜度、稳定性和成本。在 BigQuery 中,使用预定查询或带有固定窗口的 Dataflow。
- 流处理:针对无界数据的毫秒到秒级延迟处理。使用 Pub/Sub + Dataflow。通过事件时间窗口和水印处理延迟/乱序事件;确保幂等性以防止重复数据。
- 事件驱动:由变更触发(GCS 文件落定、Pub/Sub 消息)。使用 Cloud Functions 或 Cloud Run 进行无状态响应,使用 Dataflow 进行有状态处理。权衡:单个事件的处理成本与吞吐量。
- Lambda 模式:为保证准确性和支持重新处理,同时维护流处理和批处理两条路径。复杂度加倍;可考虑类似 Kappa 的简化方案,即所有内容都可以从不可变日志(例如从 Pub/Sub 归档到 Cloud Storage)中重放。
用于处理延迟数据的 Dataflow 流处理配置简例:
events
.apply(Window.into(FixedWindows.of(Duration.standardMinutes(5)))
.withAllowedLateness(Duration.standardMinutes(10))
.accumulatingFiredPanes());
领域所有权、数据产品和契约
- 所有权和 SLO:
- 每个领域团队定义并运营其数据产品,并为其设定可用性、延迟和数据质量方面的 SLO。通过 Dataplex 目录发布 SLO,并使用 Cloud Monitoring SLI(例如,分区按时完成率)进行监控。
- 契约和互操作性:
- 通过 Pub/Sub Schema Registry (Avro/Proto) 和 BigQuery 表模式来强制执行模式。对于 CSV 提取,在 Dataflow 中进行验证,并将格式错误的行路由到死信表以供分类处理。当多个引擎必须读取相同数据时,通过 Cloud Storage 中的开放格式和 BigQuery 外部表实现互操作。
- 模式演进:
- 倾向于向后兼容的变更:添加可为空的列,在 Avro/Proto 中添加可选字段,避免在没有弃用窗口期的情况下重命名/删除字段。通过版本化的契约和弃用时间表来沟通变更。
- BigQuery 示例(向后兼容的列添加):
ALTER TABLE sales.orders
ADD COLUMN coupon_code STRING;
- 对消费者的影响:
- 维护模式的语义化版本;在迁移期间同时发布 v1 和 v2 版本。对于流处理,将数据路由到带版本的主题,或在数据中包含模式版本字段。在 BigQuery 中提供授权视图,以使消费者免受物理表结构变更的影响。
- 治理和数据血缘:
- 使用 Dataplex 和 Data Catalog 进行元数据、标签(例如,PII)和数据血缘管理。在 BigQuery 中应用行级和列级安全性。为防止数据丢失,在数据提取(例如,通过 Cloud Run 或 Dataflow 转换)环节集成 Cloud DLP,以便在存储前对敏感字段进行令牌化或脱敏处理。
非功能性权衡与部署拓扑
- 可扩展性:
- BigQuery 可为分析任务弹性扩展;Bigtable 随节点数量线性扩展,但需要精心设计行键 (row-key)(例如,使用哈希或轮换前缀)以避免热点 (hotspotting)。Dataflow 自动扩缩容可应对积压任务;通过利用 Pub/Sub 的流控 (flow control) 机制来设计背压 (backpressure)。
- 延迟:
- 流式传输到 BigQuery 可提供低延迟插入,但查询时可能会有轻微延迟;设计查询时应考虑新鲜度缓冲或基于水印 (watermark) 的窗口。对于大规模的亚百毫秒 (sub-100 ms) 读取,应预先计算结果并从 Bigtable 或 Memorystore 提供服务。
- 可用性与一致性:
- Cloud Spanner 提供强一致性的全球分布式 SQL。Bigtable 提供高可用性,但在集群间是最终一致性。BigQuery 的可用性是区域性或多区域性的;为增强弹性,应将关键数据集物化到多区域。
- 成本:
- 通过分区和聚类来优化 BigQuery,以减少扫描的字节数。对于通过网络受限链路传输的小文件,进行批处理或捆绑以减少 RPC 开销。在适当的情况下,使用 BigQuery BI Engine 为缓存的交互式仪表板提供支持。
- 区域、多区域、混合云和多云:
- 区域性设计可降低延迟和成本;多区域存储(例如,BigQuery 美国/欧盟多区域、Cloud Storage 双区域/多区域)可提高持久性和位置选择性。对于灾难恢复 (DR),需定义 RPO/RTO 并复制关键数据集。在混合云场景中,使用 Datastream 进行变更数据捕获 (CDC),并使用 Transfer Appliances 或 Storage Transfer Service 进行批量迁移。对于多云场景,应在 Cloud Storage 中标准化使用开放格式,并使用可移植的计算引擎(Apache Beam/Dataflow、Dataproc 上的 Spark),同时要认识到出口流量 (egress) 和运维开销。
分层、生命周期与服务选择
- 层次分离:
- 存储层:Cloud Storage 用于原始/青铜(bronze)数据和归档;BigQuery 用于精选/服务层分析;Bigtable 用于低延迟键访问;Spanner/Cloud SQL 用于 OLTP。
- 计算层:Dataflow 用于无服务器流式/批量处理;Dataproc 用于 Spark/Hadoop 生态系统;BigQuery 用于数据仓库内 ELT;Cloud Run/Functions 用于事件驱动的微服务。
- 编排层:Cloud Composer (Airflow) 或 Workflows 用于 DAG 和 API 编排;Scheduler 用于类似 cron 的触发器。
- 服务层:Bigtable 或 Spanner 用于在线读取;BigQuery 用于 BI;Looker/BI Engine 用于仪表盘;Memorystore 用于缓存。
- 数据生命周期:
- 注入:Pub/Sub 用于流数据;Storage Transfer 或 gsutil 用于文件;Data Transfer Service 用于 SaaS。验证、去重,并将不可变的原始数据存入启用了对象版本控制的 Cloud Storage。
- 处理:使用 Dataflow 或 BigQuery 将原始数据转换为白银(silver)级(已清洗、整合),然后再转换为黄金(gold)级(业务就绪的数据集市)。
- 服务:发布 BigQuery 视图/表用于分析;将特征或预测结果预计算到 Bigtable 中供 API 使用。
- 保留与归档:应用 Cloud Storage 生命周期规则将数据迁移到 Coldline/Archive 层;使用 BigQuery 的时间分区和分区过期功能来管理数据保留。在需要时启用 CMEK,并使用 VPC Service Controls 防止数据外泄。
- 基于工作负载特征的服务选择:
- 高吞吐量、宽行、低延迟的时间序列数据:Bigtable。
- 强一致性的全球 OLTP,支持 ANSI SQL:Cloud Spanner。
- 规模适中的传统关系型事务:Cloud SQL。
- PB 级分析,支持 ANSI SQL,且存储与计算分离:BigQuery。
- 实时注入和处理:Pub/Sub + Dataflow。
- 批量 Spark/Hadoop 或需要特定库的工具:Dataproc。
BigQuery 分区简例:
CREATE TABLE ops.events
PARTITION BY DATE(event_ts)
CLUSTER BY device_id AS
SELECT * FROM staging.events_clean;
实际问题场景
Contoso Mobility 运营着一个全球性的电动滑板车车队,需要对骑行遥测和计费数据进行实时注入、处理、存储和分析。他们必须支持每分钟数百万个事件、亚秒级的欺诈规则、实时更新的仪表盘、隐私控制以及弹性的多区域运营。
方法:
- 使用 Cloud Pub/Sub 建立事件注入管道。
- 理由:Pub/Sub 提供单一的全球端点、持久化缓冲和水平扩展能力,以应对突发的设备流量。每个滑板车使用有序键,以在 1 小时窗口内保持设备内部的事件顺序。
- 使用 Cloud Dataflow (Apache Beam) 实现流式处理。
- 理由:Dataflow 的自动扩缩容可以处理流量尖峰,并与幂等键结合使用时,可为接收器(sink)提供精确一次性(exactly-once)的保证。使用事件时间窗口和水位线(watermark)来处理延迟/乱序的遥测数据。将主输出发送到精选数据流,并将一个旁路输出用于死信记录。
- 配置:
.withAllowedLateness(Duration.standardMinutes(15))
.discardingFiredPanes();
- 分别将原始数据和精选数据持久化到 Cloud Storage 和 BigQuery 中。
- 理由:将原始(bronze)的 Avro 文件存入一个双区域的 Cloud Storage 存储桶,用于数据重放和审计。将精选(silver)的流数据写入 BigQuery 的分区表进行分析,并按 scooter_id 进行聚类以实现高效的点查找。在仪表盘查询上应用一个小的刷新缓冲区,以避免短暂的流式数据陈旧问题。
- 通过 Cloud Bigtable 提供运营查询和欺诈检查服务。
- 理由:低于 100 毫秒的规则评估需要低延迟的随机访问。在 Dataflow 中预计算聚合值(例如,每台设备每 5 分钟窗口内的骑行次数),并使用哈希前缀行键(例如,h(prefix)+device_id+window_start)将其写入 Bigtable,以避免热点问题并跨 tablet 并行化读取。
- 在 Cloud Spanner 中管理事务性计费。
- 理由:计费需要全局一致的 SQL、强一致性和高可用性。在主地理区域使用一个领导者(leader)副本,并在次要区域设置只读副本,以降低客户门户的读取延迟。
- 使用 Dataplex、Data Catalog 和 Cloud DLP 实施治理。
- 理由:对 PII 字段进行分类,为数据集添加标签,并在 BigQuery 中应用列级安全性。在 Dataflow 管道中集成 Cloud DLP,在存储前对敏感属性进行令牌化。Dataplex 域反映了组织所有权;每个域都发布带有 SLO 的、文档化的数据产品。
- 使用 Cloud Composer 和 Cloud Monitoring 进行编排和运营。
- 理由:Composer 协调批量回填、数据压缩和机器学习特征物化等任务。Monitoring 监控端到端的 SLI 指标:Pub/Sub 积压量、Dataflow 水位线延迟、BigQuery 分区完整性和 Bigtable 尾部延迟。在违反 SLO 时发出警报;根据积压量增长自动扩缩容 Dataflow。
- 通过分区和分层优化成本和生命周期。
- 理由:BigQuery 表按 event_ts 进行分区,保留期为 90 天,并按 scooter_id 进行聚类。Cloud Storage 使用生命周期规则在 30 天后将原始数据迁移到 Coldline,180 天后迁移到 Archive。通过调度的 BigQuery 作业将小的微批次文件压缩成更大的 parquet 对象,以减少下游 Spark 作业的文件数量开销。
- 验证风险和弹性。
- 理由:在预期峰值 2 倍的负载下进行压力测试,以验证 Pub/Sub 配额和 Dataflow 的自动扩缩容能力。执行一次区域故障切换演练:BigQuery 多区域数据集和双区域存储桶可保持可用性;Spanner 的多区域实例通过自动故障切换维持 RPO=0 和配置的 RTO。使用基础设施即代码 (Terraform) 并结合策略验证来强制执行 CMEK 和 VPC Service Controls。
该架构清晰地分离了关注点:Pub/Sub 缓冲注入流量,Dataflow 进行计算,Cloud Storage 和 BigQuery 存储并服务于分析,Bigtable 加速运营读取,而 Spanner 保证一致性事务。它通过分区、聚类、生命周期策略和自动扩缩容,在可扩展性和延迟之间取得了平衡,并通过文档化的数据产品、契约和持续验证,将治理和可靠性融入其中。
所有领域 · 数据存储、数据湖与文件格式 →
练习这些题目 → · 在 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.
通过考试 →