Google PDE: 工作流编排与流水线自动化 — 学习指南

属于 Google Professional Data Engineer — 学习指南. 使用经过验证的答案练习: Google 考试中心, 或参加限时模拟考试: ExamRoll.io.

概览

工作流编排与流水线自动化旨在协调跨服务的数据任务,以便可靠、安全且经济高效地完成数据注入、转换、质量检查和发布。在 Google Cloud 中,编排必须与每个工作负载的执行模型保持一致:计划批处理、事件驱动的流处理、即席查询或长时间运行的作业。其设计目标是可重复性、幂等性、可观测性、最小权限以及在不同环境间的安全晋升。

关键选择:

运维模型强调采用有界指数退避的重试机制、超时、服务等级协议 (SLA)、追溯执行 (catchup) 和回填 (backfill)、确保安全重跑的幂等任务设计,以及带有死信捕获的稳健故障处理机制。安全性通过每个流水线使用独立的服务账号、密钥隔离、参数化以及最小权限 IAM 来强制执行。CI/CD、基础设施即代码 (IaC) 和全面的遥测技术共同构成了一套生产就绪的方法。

Google Cloud 上的编排:工具与模式

Cloud Composer (Airflow)

undefined

Cloud Workflows、Cloud Scheduler、Cloud Run 作业和事件驱动执行

Dataform:用于 BigQuery 的 SQL 工作流

undefined

Dataproc、Dataflow 和存储模式

可靠性、故障处理与幂等性

重试、超时与退避

undefined

和全局的 DAG SLA;在 Workflows 中,设置每个步骤的超时和带有

undefined

undefined

的重试策略。对于 Cloud Run 作业,设置重试次数和退避策略。

回填、追赶与故障处理

幂等任务设计与重跑

undefined

undefined

;使用

undefined

为流式插入去重。对于批处理,先写入一个暂存表,然后在一个事务安全的步骤中

undefined

到目标表,以允许完全重跑。

undefined

或日期对输出进行分区,写入完成标记,并使下游任务依赖于这些标记。仅重新处理标记为未完成的分区。

故障排查与可伸缩性

安全性、参数化、环境与 CI/CD

参数化与配置管理

密钥、服务账号与最小权限

CI/CD 与基础设施即代码

可观测性、告警和运行手册

遥测和告警

运行手册设计

实际问题场景

Acme 零售分析公司需要注入每日合作伙伴提供的 CSV 文件,这些文件偶尔包含格式错误的数据行。他们需要转换有效数据并加载到 BigQuery,同时将坏数据行呈现出来以供调查。他们还希望通过事件驱动的方式进行数据扩充,以实现近实时的价格更新,并确保从开发环境到生产环境的安全提升。

方法:

  1. 存储和事件触发器

    • 创建一个专用的 Cloud Storage 存储桶,启用对象版本控制和统一的存储桶级访问权限。通过 Eventarc 将对象完成通知发送到 Pub/Sub。
    • 理由:对象完成是一个可靠的事件,可用于触发下游注入;版本控制支持重新运行和审计。
  2. 带死信处理的批量注入

    • 使用 Cloud Composer 调度一个每日 02:00 运行的 Airflow DAG,并启用追溯 (catchup)。该 DAG 启动一个 Dataflow 批量作业,该作业解析 CSV、验证模式,并将有效记录写入 BigQuery 的确定性暂存表,然后通过 MERGE 操作合并到分区目标表中。将格式错误/失败的记录路由到 BigQuery 的死信表。
    • 理由:Dataflow 可扩展解析/验证过程;MERGE 确保幂等性;死信捕获支持在不阻塞流水线的情况下进行检查,这符合针对格式错误行的推荐模式。
  3. 事件驱动的扩充

    • 部署一个 Cloud Run 作业,为增量价格更新执行轻量级扩充。当白天有小型更新文件到达时,通过 Cloud Workflows 监听来自 Eventarc 的 Pub/Sub 消息来触发该作业。
    • 理由:使用 Workflows 的无服务器容器为小型事件提供低延迟、低运维的编排,同时将繁重的转换保留在批处理中。
  4. 可靠性控制

    • 为 Dataflow 和 Cloud Run 作业中的瞬时故障配置指数退避重试,并将总重试时间限制在 DAG SLA 内。在 Airflow 中设置每个任务的执行超时和 on_failure 回调;在 Workflows 中,设置 max_doublings 和 max_retry_duration。
    • 理由:有界的退避可保护 SLA 并防止失控的重试。
  5. 安全和最小权限

    • 每个组件都使用专用的服务账号运行:Composer 编排器 SA、Dataflow 工作器 SA、Cloud Run 作业 SA。仅授予必需的角色:授予 Dataflow 对注入存储桶的 GCS 读取权限,对目标数据集的 BigQuery dataEditor 权限,以及对日志的 Viewer 权限。将密钥存储在 Secret Manager 中,并在运行时引用它们。
    • 理由:强制执行最小权限原则并隔离爆炸半径。
  6. 元数据驱动的编排

    • 维护一个 BigQuery 控制表,其中列出合作伙伴来源、文件模式和目标数据集。在 DAG 运行时,Airflow 查询此表并使用动态任务映射为每个合作伙伴生成任务。
    • 理由:添加合作伙伴变成数据变更,而非代码变更,从而降低部署风险。
  7. 可观测性和告警

    • 发出带有 run_id 和 partner_id 的结构化日志。为 DAG SLA 未达成、Dataflow 系统延迟和死信计数不为空等情况创建告警策略。对于向目标表插入 BigQuery 数据的操作,配置一个 Cloud Logging 接收器,使用高级过滤器筛选该表的日志,并将其发送到 Acme 监控工具所消费的 Pub/Sub 主题。
    • 理由:细粒度的告警能够实现快速分类,且无噪声干扰。
  8. CI/CD 和提升

    • 在 Terraform 中管理基础设施(存储桶、Pub/Sub、Eventarc、Composer、Workflows、BigQuery 数据集)。使用 Cloud Build 验证 Airflow DAG 语法、运行单元测试,并部署到开发环境的 Composer。在 Dataform 断言和集成测试通过后,使用参数化配置和手动批准门控将其提升到测试和生产环境。
    • 理由:声明式、可重复的部署以及跨环境的安全提升。
  9. 运行手册和恢复

    • 记录重放特定日期数据的步骤:从对象版本控制中恢复 CSV,为该分区重新运行 Dataflow 作业,MERGE 结果,并审查 DLQ 记录。包括一个“fixed dataset replay”流程,以便在出现差异时隔离转换逻辑中的错误。
    • 理由:幂等设计和文档化的恢复流程简化了部分故障的修复工作。

数据注入、集成与迁移 · 所有领域 · 机器学习、AI 与数据服务

练习这些题目 → · 在 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.

通过考试 →

浏览 Google →

Related guides

一体化访问

一次订阅。所有考试。

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

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

无需信用卡*

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

无需信用卡*

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