Google PDE: 数据注入、集成与迁移 — 学习指南
属于 Google Professional Data Engineer — 学习指南. 使用经过验证的答案练习: Google 考试中心, 或参加限时模拟考试: ExamRoll.io.
概览
Google Cloud 中的数据注入、集成和迁移涵盖了多种可重复的模式、托管服务和运维控制,可将各种源系统转变为可靠、可查询的数据集。有效的设计应将传输与转换分离,解耦生产者和消费者,并倾向于采用具有清晰沿袭和验证机制的、幂等的、带检查点的流水线。本节内容涵盖注入模式、用于数据移动和 CDC 的 Google Cloud 服务、模式和数据质量控制、连接性和混合集成以及切换策略,并会贯穿全文指出其中的设计权衡和故障模式。
注入模式和工作负载
- 批量注入:按定义的间隔定期拉取或放置文件。适用于可预测的成本和数据回填。故障模式:大批量、低频率的批处理会导致资源尖峰、过长的追赶窗口和错过 SLA。缓解措施:适当调整批处理窗口大小,按时间或键进行分片,并使用并行处理。
- 批量加载:一次性或大规模加载(例如,初始历史数据回填)。首选列式或自描述格式(Parquet、Avro),并直接加载到分析存储(BigQuery)或暂存在 Cloud Storage 中。权衡:通过外部表查询可避免加载步骤,但会将成本转移到查询时的扫描上。
- 增量加载:通过时间戳或高水位标记定期进行增量加载。需要稳健的去重机制和幂等的更新插入 (upsert) 操作。故障模式:时钟偏斜或记录迟到。应使用服务器端提交时间戳和水印 (watermarking)。
- 变更数据捕获 (CDC):持续复制来自业务数据库的插入、更新和删除操作。最适合近实时分析和低停机时间迁移。权衡:
- 顺序:大多数 CDC 工具会保留事务内的顺序,并且通常在分片内也能保持顺序,但不保证跨分片的全局顺序。应使用事务提交时间戳和主键来重构序列。
- 交付语义:通常是“至少一次” (at-least-once);应构建幂等的接收器或使用唯一的变更 ID 进行去重。
- 快照 + CDC:从一个一致的快照开始,然后应用来自精确日志序列的变更,以在无停机的情况下达到数据同步。
关系型、SaaS、本地和文件源:
- 关系型源:使用原生 CDC 或时间戳列。对于批量加载,导出为 Avro/Parquet 并暂存在 Cloud Storage 中。
- SaaS 源:首选带有增量令牌的供应商 API;通过托管连接器(例如,在 Data Fusion 中)进行集成。注意速率限制并处理模式漂移。
- 本地源:可选择基于代理的传输、VPN/Interconnect + Private Google Access,或使用 Transfer Appliance 进行离线植入。
- 文件注入:对于大量小文件,应进行打包(例如,tar)以减少 RPC 开销。使用
undefined
或并行化客户端;对于分析任务,应将文件组合或转换为更大的列式文件。
用于注入、集成和迁移的 Google Cloud 服务
- Datastream (无服务器 CDC):从 MySQL、PostgreSQL 和 Oracle 捕获变更到 Cloud Storage、BigQuery(通过模板)或 Pub/Sub。它保留了事务边界和提交元数据;不保证全局顺序。应在下游按键和提交时间戳进行排序。交付语义为“至少一次” (at-least-once);需设计幂等的消费者(例如,使用变更 ID 的 BigQuery MERGE 操作)。
- Database Migration Service (DMS):用于使用原生复制实现停机时间最短的数据库迁移。DMS 会创建一个一致性快照,然后使用 GTID/LSN/SCN 持续复制变更。它专为直接迁移 (lift-and-shift) 而构建,不适用于任意转换。对于分析场景,如果需要,可以使用 Dataflow 或 Data Fusion 来增强 DMS。
- Cloud Data Fusion:一种托管集成服务,提供与关系型数据库、SaaS、文件和消息传递系统的连接器。可以用它构建包含转换阶段(连接、聚合、格式转换、自定义 Wrangler 配方)的流水线,并捕获跨源和字段的沿袭。在运维方面,它能进行调度、重试和发出指标。可使用 Data Fusion 进行无代码/低代码的 ELT/ETL,并集中管理连接器。
- Storage Transfer Service (STS):托管的、可调度的传输服务,可将数据从 AWS S3、Azure Blob、本地(使用代理)、SFTP 和 URL 列表传输到 Cloud Storage。支持清单、增量同步、带宽控制和校验和完整性。故障模式包括小文件效率低下和 API 节流;可通过批处理和调整并发度来缓解。
- Transfer Appliance:一种离线加密设备,适用于当网络带宽有限或数据过于敏感不宜长时间传输时,进行 TB 到 PB 规模的初始数据植入。内置了监管链和加密功能。植入数据后,可接着使用 STS 或 CDC 处理增量数据。
- Cloud Pub/Sub + Dataflow:Pub/Sub 为流式或微批处理模式解耦了生产者和消费者。Dataflow 提供自动扩缩、有状态的流/批处理能力,并带有检查点和水印功能。使用 BigQuery Storage Write API 可实现低延迟流式传输,并为每个默认流提供精确一次性 (exactly-once) 保证;否则,需依赖 insertId 的去重语义。
对于 Hadoop 到 Dataproc 的迁移,应通过 GCS 连接器将数据存储在 Cloud Storage 中,并使用临时或自动扩缩集群,从而最大限度地减少对 Persistent Disk 的使用。这样可以避免高昂的块存储成本,同时为处理过程保留与 HDFS 兼容的语义。
边界处的模式、验证与数据质量
- 模式映射与类型转换:尽早标准化为强类型模式。Avro 或 Parquet 能够保留模式并支持平滑演进。在 BigQuery 中,优先使用分区表和集群表以降低扫描成本。 示例:创建一个用于日常分析的分区表 CREATE TABLE dataset.tracking_table ( event_ts TIMESTAMP, device_id STRING, payload STRING ) PARTITION BY DATE(event_ts) CLUSTER BY device_id;
- 格式错误记录的处理:将拒绝的记录路由到死信队列 (Pub/Sub) 或 Cloud Storage 中的隔离存储桶。在 Dataflow 中使用旁路输出,或在 Data Fusion 中使用错误收集器。记录解析错误,并附上示例负载和模式版本,以便进行分类处理。
- 验证:在持久化之前执行边界检查:
- 结构性:模式一致性、必填字段、数据类型、枚举域。
- 引用性:通过缓存的维度查找来验证外键是否存在。
- 合理性:时间戳的范围、地理围栏、非负金额。
- 唯一性:主键或复合键冲突。
- 幂等加载:使用确定性键和 upsert (更新或插入) 操作。在 BigQuery 中,使用自然键或代理变更键来实现 MERGE 操作。 示例: 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);
- 水印与延迟数据处理:在流处理管道中,配置事件时间水印和允许的延迟时间,以平衡完整性和延迟。延迟数据被路由到修正路径或触发回填。
- 对账:跟踪从源到接收器每个分区/窗口的行数和校验和。捕获 CDC 日志位置 (LSN/SCN) 和提交时间戳;将其存储在控制表中以证明连续性并识别数据间隙。
连接性、可靠性与运维
网络连接与私有访问:
- 混合云:使用 Cloud VPN 或 Dedicated/Partner Interconnect 建立私有连接。启用 Private Google Access 或 Private Service Connect 以私密访问 Cloud Storage 等 Google API。
- 安全性:使用服务账号作为工作负载身份,遵循最小权限原则的 IAM,使用 VPC Service Controls 防止数据泄露,并在需要时使用 CMEK。
- 吞吐量:在客户端扩展并行度,但最终带宽决定吞吐量。对于大规模传输,首选 Transfer Appliance 进行初始批量传输,然后使用 STS 或 CDC 进行增量更新。
检查点与反压:Dataflow 负责管理检查点和自动扩缩;设计的接收器应能吸收突发流量 (例如缓冲到 Cloud Storage,批量写入 BigQuery)。对于 Pub/Sub,调整流控制和确认截止时间,以防止消息重传风暴。
使用 CDC 实现有序性与一致性:
- Datastream 保留事务内顺序并发出提交元数据;消费者使用提交时间戳来重构每个键的顺序。预期为“至少一次”交付,因此需要构建幂等性处理。
- DMS 使用原生日志确保快照和复制切换过程中的数据库一致性。使用只读副本或双写策略进行分阶段切换。
分析场景的文件策略:对于需要多引擎访问的大型数据集,将规范数据存储在 Cloud Storage 中,并在成本效益允许的情况下,通过永久外部表支持即席查询。对于生产分析,将数据加载到 BigQuery 分区表中,以最小化每次查询的扫描成本。
小文件优化:在传输前将小文件打包 (例如,每个 tar 包约 1000 个文件),然后在云端解压。使用并行化的 gsutil 和生命周期规则对暂存产物进行分层和过期处理。
运维陷阱与缓解措施:
- 来自 SaaS 的模式漂移:在 Data Fusion 中启用模式演进并强制执行兼容性。对破坏性变更进行告警。
- 时区与编码:在入口处统一规范化为 UTC 和 UTF-8。
- CDC 中的数据间隙:监控源端日志保留期;当副本延迟接近保留期限制时发出告警。
- 配额:BigQuery 流式插入、API 速率限制;在接近限制时进行批处理。
切换、回填与验证
- 切换规划:
- 大爆炸式 (Big bang):冻结期短,一次性切换。操作复杂度最低;但如果需要回滚,风险最高。
- 分阶段或蓝绿部署 (Phased or blue/green):双轨运行,镜像写入,逐步转移流量,并进行影子读取。成本较高;但回滚更安全。
- 回填:
- 执行初始批量加载(使用 Transfer Appliance 或 STS),采用 Avro/Parquet 格式以保留 schema。在加载过程中进行分区和聚类,以避免返工。
- 在进行快照的同时,从一个已知的日志位置启动 CDC,以捕获批量传输期间的增量变化。在向生产环境开放前,在一个共同的水位线 (watermark) 进行对账。
- 回滚:
- 在验证期间,保持旧系统为只读状态。对于双写场景,将写入操作置于功能标志 (feature flag) 之后,以便快速恢复。保留一个一致的检查点,以便在需要时重放或撤销 CDC 变更。
- 迁移验证:
- 结构性验证:行数和各分区的校验和 (checksum) 匹配;schema 和约束等效。
- 时间性验证:从快照边界到切换点无数据间隙;CDC 位置连续。
- 业务对等性验证:在时间窗口内比较聚合指标和 KPI;运行验收查询。
- 性能验证:验证摄取吞吐量、查询延迟和成本是否符合预算。
实践问题场景
Northstar Retail 公司必须将其全球混合部署的本地 Oracle 和 MySQL 事务系统、SaaS CRM 事件以及每日的 CSV 文件整合到 Google Cloud 中,以支持近乎实时的分析和机器学习。他们还需要迁移一个旧有的 Hadoop 集群,同时避免产生高昂的块存储费用,并实现零或低停机时间切换。
- 建立安全的混合连接
- 使用 Partner Interconnect 作为主带宽,Cloud VPN 作为备用。启用 Private Google Access,以便本地工作负载可以私密地访问 Cloud Storage 和 Pub/Sub。 理由:私有路径最大限度地减少了出口暴露和延迟,而 Private Google Access 在满足安全策略的同时,避免了对公网 IP 的需求。
- 高效地植入历史数据
- 对于 800 TB 的历史 HDFS 数据,使用 Transfer Appliance(初始批量)复制到 Cloud Storage。植入数据后,每天从本地 NFS 导出运行 Storage Transfer Service,以在切换前获取变更。 理由:Transfer Appliance 避免了长时间的网络饱和;STS 提供定期的、带校验和的增量同步。通过 GCS connector 将数据存储在 Cloud Storage 中,允许 Dataproc 在无需每个节点挂载 50 TB Persistent Disk 的情况下处理数据。
- 使用 CDC 迁移操作型数据库
- 使用 DMS 以最小的停机时间迁移 MySQL 和 PostgreSQL。对于 Oracle 到分析系统的 CDC,使用 Datastream 将数据落地到 Cloud Storage,然后使用 Google 提供的 Dataflow 模板加载到 BigQuery。 理由:DMS 利用原生复制功能实现可靠的快照 + 持续同步;Datastream 提供带有提交元数据的无服务器 CDC,而 Dataflow 模板确保了有序、幂等的 BigQuery 写入。
- 摄取 SaaS 和基于文件的供给
- 使用 SaaS 连接器构建 Cloud Data Fusion 管道,以处理带有增量令牌的 CRM 事件,并构建一个文件管道,通过 STS 从供应商的 SFTP 服务器摄取每日 CSV 文件。将数据规范化为 Avro 格式,存入一个策展后的 Cloud Storage 存储桶,然后加载到分区的 BigQuery 表中。 理由:Cloud Data Fusion 集中管理连接器、转换和数据血缘。标准化为 Avro 格式保留了 schema 并简化了演进;分区的 BigQuery 表降低了查询成本。
- 流式处理实时事件
- 将网站和商店事件发布到 Pub/Sub。使用 Dataflow 进行解析、验证、丰富和水印处理;通过 Storage Write API 写入 BigQuery,并将原始 Avro 归档到 Cloud Storage。 理由:Pub/Sub 解耦了生产者/消费者;Dataflow 提供自动扩缩、有状态处理、检查点和延迟数据处理;双写确保了低延迟分析和持久的原始数据保留。
- 实施边界数据质量和 schema 控制
- 在 Dataflow/Data Fusion 中实施 schema 注册和验证。将格式错误的数据路由到一个 GCS 隔离存储桶和一个 Pub/Sub 死信主题。应用领域检查(例如,货币代码、UTC 时间戳)并使用复合键进行去重。 理由:早期拒绝和隔离可防止坏数据扩散;幂等性和去重可防范来自 CDC 和流式处理源的“至少一次”交付所带来的问题。
- 优化分析存储和访问
- 将策展后的数据集加载到分区和聚类的 BigQuery 表中。将原始归档数据作为永久外部表暴露,用于低频探索。对于仍保持事务性的 OLTP 工作负载,保留带有只读副本的 Cloud SQL。 理由:分区和聚类最大限度地减少了扫描成本;外部表避免了为偶尔访问而进行不必要的数据加载;Cloud SQL 为事务性应用保留了 ACID 语义。
- 规划切换、回填和回滚
- 为每个 RDBMS 执行快照 + CDC;达到一个对账点,在该点行数和校验和匹配。进行 48 小时的蓝绿部署和双写,逐步将读取流量转移到 BigQuery。维护一个功能标志,以便在检测到差异时恢复写入操作。 理由:蓝绿部署降低了风险;在已知水位线进行验证确保了完整性;功能标志实现了快速回滚。
- 验证与可观测性
- 构建控制表,捕获每个分区的源 LSN/SCN、提交时间戳、行数和校验和。监控 Datastream 延迟、DMS 复制状态、Dataflow 水位线、Pub/Sub 积压、STS 作业状态以及 BigQuery 流式插入指标。 理由:端到端的数据血缘和量化控制为正确性提供了可审计的证明,并能及时对数据间隙或延迟发出警报。
通过分离落地层、策展层和服务层;使用 Cloud Storage 作为持久、低成本的暂存和归档区;利用 DMS/Datastream 实现 CDC 并配合幂等消费者;以及在入口处强制执行 schema 和质量控制,Northstar Retail 实现了一个安全、可扩展的摄取方案,并以可预测的成本完成了一次低风险、可验证的迁移。
← Spark、Dataproc 与分布式数据处理 · 所有领域 · 工作流编排与流水线自动化 →
练习这些题目 → · 在 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.
通过考试 →