Google PDE: 工作流编排与流水线自动化 — 学习指南
属于 Google Professional Data Engineer — 学习指南. 使用经过验证的答案练习: Google 考试中心, 或参加限时模拟考试: ExamRoll.io.
概览
工作流编排与流水线自动化旨在协调跨服务的数据任务,以便可靠、安全且经济高效地完成数据注入、转换、质量检查和发布。在 Google Cloud 中,编排必须与每个工作负载的执行模型保持一致:计划批处理、事件驱动的流处理、即席查询或长时间运行的作业。其设计目标是可重复性、幂等性、可观测性、最小权限以及在不同环境间的安全晋升。
关键选择:
- 使用 Cloud Composer (Apache Airflow) 进行以代码为中心的批处理编排,适用于 DAG、任务依赖和高级调度。
- 使用 Cloud Workflows 进行无服务器 API 协同,适用于轻量级、事件驱动的跨服务序列。
- 使用 Cloud Run 作业或 Dataproc 作业等执行端点,由 Cloud Scheduler(用于 cron 作业)或 Eventarc(用于事件)触发。
- 使用 Dataform 进行原生 SQL 编排,适用于 BigQuery 转换、断言和发布管理。
运维模型强调采用有界指数退避的重试机制、超时、服务等级协议 (SLA)、追溯执行 (catchup) 和回填 (backfill)、确保安全重跑的幂等任务设计,以及带有死信捕获的稳健故障处理机制。安全性通过每个流水线使用独立的服务账号、密钥隔离、参数化以及最小权限 IAM 来强制执行。CI/CD、基础设施即代码 (IaC) 和全面的遥测技术共同构成了一套生产就绪的方法。
Google Cloud 上的编排:工具与模式
Cloud Composer (Airflow)
- DAG 定义了具有明确依赖关系的有向无环执行图。使用 TaskFlow API 或算子(例如,BigQuery、Dataflow、Dataproc、Cloud Run)来表达任务。传感器 (Sensor) 和可延迟算子 (deferrable operator) 可减少调度器在等待条件(例如,Cloud Storage 中的对象完成事件或 BigQuery 中出现新分区)下的负载。
- 调度:cron 表达式、start_date、end_date 和 catchup 控制历史运行。使用 catchup 进行回填;对于邻近流处理或非幂等的目标则禁用它。使用 max_active_runs 和池 (pool) 限制并发,以保护下游系统。
- 依赖关系:使用 set_upstream/set_downstream 或 taskflow 依赖。对于元数据驱动的编排,可使用动态任务映射从 BigQuery 控制表(例如,客户/分区列表)动态生成任务,从而保持 DAG 解析时间稳定并使任务由数据驱动。
- 示例(简洁的)DAG 片段:
undefined
Cloud Workflows、Cloud Scheduler、Cloud Run 作业和事件驱动执行
- Cloud Workflows 编排 Google API 和 HTTP 端点,并内置了重试、循环、并行分支和补偿逻辑。它非常适合跨服务(如 BigQuery、Dataflow、Batch 和 Cloud Run 作业)的轻量级控制流。
- Cloud Scheduler 触发 Workflows、Pub/Sub 主题或 HTTP 服务,以实现 cron 风格的自动化。对于每日凌晨 02:00 的批处理,可以调度一个 Workflow 来启动 Dataflow 作业或 Dataproc 作业。
- Cloud Run 作业执行容器化的批处理步骤,具有自动重试和最少的运维操作。它们与 Workflows 搭配使用,非常适合多步骤数据任务或在 Dataflow 或 BigQuery 周围进行预/后处理。
- 事件驱动:使用 Eventarc 将 Cloud Storage 对象完成事件、Pub/Sub 消息或 Audit Logs 路由到 Cloud Run 或 Workflows。对于单个表的 BigQuery 插入作业通知,可以创建一个带有高级过滤器的 Cloud Logging 接收器 (sink) 将日志发送到 Pub/Sub,然后从该主题触发您的消费者。
Dataform:用于 BigQuery 的 SQL 工作流
- 使用 ref() 建立模型依赖关系图,定义表/视图/增量表,并按标签或计划编排构建。Dataform 将 SQLX 编译为有序的执行计划,从而能够通过声明式定义实现元数据驱动的编排。
- 断言 (Assertion) 确保数据质量。断言是一个查询,它必须返回零行才算通过。 断言示例: – definitions/assert_non_negative_prices.sqlx
undefined
- 发布和代码库控制:将代码存储在代码库中,使用分支和审查,并使用特定于环境的变量将带有标签的发布版本晋升到不同环境(例如,dev、test、prod)。通过 CI/CD 检查和断言结果来控制部署。
Dataproc、Dataflow 和存储模式
- 为了以最少的运维操作重用 Hadoop/Spark,请将 Dataproc 与 GCS 连接器结合使用,使数据能够在集群生命周期之外持久化,并最大限度地降低永久性磁盘成本。为每个作业创建临时集群以实现隔离和成本控制;使用 Composer 或 Workflows 进行编排。
- 对于包含格式错误行的批处理注入,运行 Dataflow 将有效记录写入 BigQuery,并将解析/验证错误路由到死信 BigQuery 表中以供检查。
可靠性、故障处理与幂等性
重试、超时与退避
- 对于瞬时故障,使用有界指数退避,并将总重试窗口限制在作业的 SLA 内。例如,一个每 15 分钟轮询一次数据库的前端或任务,应使用指数退避重试最多 15 分钟,然后呈现一个受控的失败。
- 在 Airflow 中,配置每个任务的
undefined
和全局的 DAG SLA;在 Workflows 中,设置每个步骤的超时和带有
undefined
及
undefined
的重试策略。对于 Cloud Run 作业,设置重试次数和退避策略。
回填、追赶与故障处理
- 当任务是幂等的且源数据按日期分区时,启用 catchup 进行历史数据重计算。对于非确定性输出或外部副作用,考虑使用仅用于回填的 DAG 或写入审计表来跟踪已生成的内容。
- 在流式/批处理转换中,对记录级别的失败使用死信主题/表。对于批处理 Dataflow,使用错误标签捕获坏数据行并聚合错误指标;对于流处理,使用 Pub/Sub DLQ。
幂等任务设计与重跑
- BigQuery:首选使用带有去重键的
undefined
或
undefined
;使用
undefined
为流式插入去重。对于批处理,先写入一个暂存表,然后在一个事务安全的步骤中
undefined
到目标表,以允许完全重跑。
- Cloud Storage:使用 generation 前置条件和确定性的对象名称(例如,前缀/日期/哈希值),以便重跑仅在预期情况下安全地覆盖。
- Pub/Sub 和 Dataflow:为至少一次交付进行设计。包含消息标识符(例如,包裹 ID、逻辑事件时间戳),以便下游系统可以去重和判断延迟情况。如果业务规则接受“先处理的事件获胜”的语义,请记录这一权衡并监控数据倾斜;否则,通过事件时间与决胜规则来确定获胜者。
- 从部分失败中恢复:按
undefined
或日期对输出进行分区,写入完成标记,并使下游任务依赖于这些标记。仅重新处理标记为未完成的分区。
故障排查与可伸缩性
- 当流式仪表板丢失事件但 Pub/Sub 显示事件存在时,通过 Dataflow 流水线运行一个已知的固定数据集,以隔离转换逻辑中的缺陷。验证窗口、触发器和允许的延迟。
- 常见失败模式:为无界数据源创建流处理流水线时,未使用适当的窗口/触发器,或不正确地使用分片窗口,可能导致流水线创建失败或状态爆炸。
- 通过 max workers 和自动扩缩算法来扩展 Dataflow;对于流量尖峰(例如,50,000 次安装),提高最大工作器数量以允许在高峰期间进行水平扩展。
安全性、参数化、环境与 CI/CD
参数化与配置管理
- 按环境将配置外部化。在 Composer 中,使用 Variables、Connections 和环境变量;按执行日期或分区对 DAG 参数进行模板化。在 Workflows 中,使用运行时参数并为每个环境使用独立的工作流,或从 Secret Manager 读取配置。
- 通过读取一个列出客户端、数据源或分区的控制表(例如,BigQuery 配置数据集)来实现元数据驱动的编排。动态生成任务,从而将代码变更与数据驱动的变更解耦。
密钥、服务账号与最小权限
- 将凭证存储在 Secret Manager 中,并在运行时引用它们。避免将密钥硬编码在代码或 Airflow Variables 中。
- 为每个流水线分配一个独立的服务账号,并授予所需的最小 IAM 角色。对于受监管的 BigQuery 访问,将客户数据隔离到独立的数据集中,仅向经批准的用户授予特定于数据集的角色,并将 BigQuery API 访问权限限制为经批准的主体。对于多租户场景,为每个客户创建一个数据集,并仅绑定适当的角色。
CI/CD 与基础设施即代码
- 使用 Terraform 管理基础设施(Composer 环境、Workflows、Scheduler 作业、Pub/Sub 主题、日志接收器)。使用模块来标准化项目/环境、密钥和服务账号。
- 使用 Cloud Build 或 GitHub Actions 构建和测试流水线代码。自动化单元测试、SQL 语法检查、Dataform 空跑和 Airflow DAG 验证。通过标签提升构建产物;对于 Composer,将 DAG 打包为可部署的捆绑包;对于 Dataform,使用在断言通过后进行提升的发布分支。
- 部署提升:通过独立的项目和参数化配置实现 dev → test → prod。对于高风险的提升,使用带有手动审批门和变更窗口的持续交付。
可观测性、告警和运行手册
遥测和告警
- 将所有编排日志以结构化字段(pipeline、dag_id、run_id、task_id、partition)路由到 Cloud Logging。通过基于日志的指标将错误日志导出到 Monitoring。在以下情况发出告警:
- 错过调度或 SLA 未达成
- 连续任务失败
- 积压增长(例如,Pub/Sub 未确认消息、Dataflow 系统延迟)
- 数据质量断言失败
- Cloud Composer:监控 DAG/任务持续时间、成功率、队列深度和调度器健康状况。配置 on_failure_callback 以进行呼叫和执行修复运行手册。
- Cloud Workflows:检查执行日志和步骤延迟;添加显式重试和错误处理程序;使用关联 ID 发出自定义日志。
- BigQuery 表变更通知:创建一个项目级的 Logging 接收器,并使用高级过滤器筛选针对特定表的插入作业,然后导出到 Pub/Sub;您的监控工具订阅该主题以获取即时告警,且不受其他表的噪声干扰。
运行手册设计
- 为每个流水线记录触发器、依赖项、SLA、回滚/重试过程以及安全的回填步骤。包括 Dataflow 的“fixed dataset replay”、如何排空流处理作业、如何重新处理失败的分区以及如何修复 DLQ 消息。
- 捕获常见的故障特征(例如,权限被拒绝、超出配额、模式不匹配),并附上决策树和上报路径。
实际问题场景
Acme 零售分析公司需要注入每日合作伙伴提供的 CSV 文件,这些文件偶尔包含格式错误的数据行。他们需要转换有效数据并加载到 BigQuery,同时将坏数据行呈现出来以供调查。他们还希望通过事件驱动的方式进行数据扩充,以实现近实时的价格更新,并确保从开发环境到生产环境的安全提升。
方法:
存储和事件触发器
- 创建一个专用的 Cloud Storage 存储桶,启用对象版本控制和统一的存储桶级访问权限。通过 Eventarc 将对象完成通知发送到 Pub/Sub。
- 理由:对象完成是一个可靠的事件,可用于触发下游注入;版本控制支持重新运行和审计。
带死信处理的批量注入
- 使用 Cloud Composer 调度一个每日 02:00 运行的 Airflow DAG,并启用追溯 (catchup)。该 DAG 启动一个 Dataflow 批量作业,该作业解析 CSV、验证模式,并将有效记录写入 BigQuery 的确定性暂存表,然后通过 MERGE 操作合并到分区目标表中。将格式错误/失败的记录路由到 BigQuery 的死信表。
- 理由:Dataflow 可扩展解析/验证过程;MERGE 确保幂等性;死信捕获支持在不阻塞流水线的情况下进行检查,这符合针对格式错误行的推荐模式。
事件驱动的扩充
- 部署一个 Cloud Run 作业,为增量价格更新执行轻量级扩充。当白天有小型更新文件到达时,通过 Cloud Workflows 监听来自 Eventarc 的 Pub/Sub 消息来触发该作业。
- 理由:使用 Workflows 的无服务器容器为小型事件提供低延迟、低运维的编排,同时将繁重的转换保留在批处理中。
可靠性控制
- 为 Dataflow 和 Cloud Run 作业中的瞬时故障配置指数退避重试,并将总重试时间限制在 DAG SLA 内。在 Airflow 中设置每个任务的执行超时和 on_failure 回调;在 Workflows 中,设置 max_doublings 和 max_retry_duration。
- 理由:有界的退避可保护 SLA 并防止失控的重试。
安全和最小权限
- 每个组件都使用专用的服务账号运行:Composer 编排器 SA、Dataflow 工作器 SA、Cloud Run 作业 SA。仅授予必需的角色:授予 Dataflow 对注入存储桶的 GCS 读取权限,对目标数据集的 BigQuery dataEditor 权限,以及对日志的 Viewer 权限。将密钥存储在 Secret Manager 中,并在运行时引用它们。
- 理由:强制执行最小权限原则并隔离爆炸半径。
元数据驱动的编排
- 维护一个 BigQuery 控制表,其中列出合作伙伴来源、文件模式和目标数据集。在 DAG 运行时,Airflow 查询此表并使用动态任务映射为每个合作伙伴生成任务。
- 理由:添加合作伙伴变成数据变更,而非代码变更,从而降低部署风险。
可观测性和告警
- 发出带有 run_id 和 partner_id 的结构化日志。为 DAG SLA 未达成、Dataflow 系统延迟和死信计数不为空等情况创建告警策略。对于向目标表插入 BigQuery 数据的操作,配置一个 Cloud Logging 接收器,使用高级过滤器筛选该表的日志,并将其发送到 Acme 监控工具所消费的 Pub/Sub 主题。
- 理由:细粒度的告警能够实现快速分类,且无噪声干扰。
CI/CD 和提升
- 在 Terraform 中管理基础设施(存储桶、Pub/Sub、Eventarc、Composer、Workflows、BigQuery 数据集)。使用 Cloud Build 验证 Airflow DAG 语法、运行单元测试,并部署到开发环境的 Composer。在 Dataform 断言和集成测试通过后,使用参数化配置和手动批准门控将其提升到测试和生产环境。
- 理由:声明式、可重复的部署以及跨环境的安全提升。
运行手册和恢复
- 记录重放特定日期数据的步骤:从对象版本控制中恢复 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.
通过考试 →