Amazon DEA-C01: 数据编排和工作流管理 — 学习指南
属于 Amazon Data Engineer Associate DEA-C01 — 学习指南. 使用经过验证的答案练习: Amazon 考试中心, 或参加限时模拟考试: ExamRoll.io.
编排和工作流管理是构建可靠、可维护的数据平台的核心:它们协调提取-转换-加载 (ETL) 作业、管理依赖关系、处理故障并集成事件驱动的流程。该领域涵盖了用于批量 ETL、复杂 DAG、无服务器状态机和事件调度的托管式 AWS 选项——每种选项都有不同的执行语义、持久性和扩展性权衡。理解何时使用 AWS Glue Workflows、MWAA、Step Functions 或 EventBridge Scheduler,以及如何配置错误处理和可观测性,对于实现可预测的管道和控制运营成本至关重要。
AWS Glue Workflows 和触发器
AWS Glue Workflows 将 Glue 作业、爬网程序和触发器组合成一个依赖关系图,让您能够运行协调的 ETL。通过控制台或 CLI (aws glue create-workflow –name MyWorkflow) 创建工作流。触发器附加到工作流上,有三种类型:计划、按需和条件。为计划触发器创建的 CLI 示例:
- aws glue create-trigger –name hourly-trigger –workflow-name MyWorkflow –type SCHEDULED –schedule “cron(0 * * * ? *)” –actions ‘[{“JobName”:“etl-job”}]’
条件触发器使用一个引用作业名称和状态 (SUCCEEDED, FAILED) 的断言 (Predicate)。断言 JSON 示例:{“Logical”:“AND”,“Conditions”:[{“JobName”:“prev-job”,“State”:“SUCCEEDED”}]}。默认情况下,Glue 条件触发器在成功时触发;要处理失败,请使用 State=FAILED 配置条件,或创建一个显式的 FAILED 触发器,将错误路由到修复作业或 SNS 警报。
操作模式和决策标准:
- 当您需要原生的 Glue 作业/爬网程序编排和血缘关系时,请使用 Glue Workflows;选择触发器以进行 cron 调度或在作业完成后进行链式调用。
- 对于临时调用,请使用 aws glue start-workflow-run –name MyWorkflow,或对按需触发器使用 start-trigger。
- 对于复杂的分支或非 Glue 任务,首选 Step Functions 或 MWAA;当管道以 Glue 为中心时,Glue Workflows 是最佳选择。
错误处理:添加 FAILED 触发器,为作业成功/失败发出 CloudWatch 指标,并通过 Lambda 将失败推送到 SQS/SNS 死信队列,以进行自动重试和调查。
用于复杂 DAG 的 Amazon MWAA (Managed Airflow)
MWAA 提供了一个托管的 Apache Airflow 环境,用于表达复杂的 DAG、任务依赖、传感器 (sensor) 和自定义算子 (operator)。使用 aws mwaa create-environment –name MyEnv –airflow-configuration-options Key=core.executor,Value=CeleryExecutor 创建环境,并提供 DAG 的 S3 路径和执行角色。重要的规模和网络细节:
- MWAA 需要一个包含私有子网和 NAT 网关的 VPC 以访问互联网;不支持仅使用公有子网的设置。
- 工作节点 (worker) 和调度器 (scheduler) 的行为通过在创建环境时提供的 Airflow 配置选项 (AirflowConfigurationOptions) 进行控制。调整 celery.worker_concurrency、celery.worker_autoscale 和调度器设置,以匹配任务并发性和 DAG 复杂性。
- 监控 CloudWatch 指标 (SchedulerHeartbeat, TasksFailed, TasksRunning, QueuedTasks),并在观察到队列增长时扩展工作节点的自动伸缩规模或增加最大工作节点数。
决策标准:
- 当您需要 Airflow 的功能时,请使用 MWAA:复杂的 DAG、丰富的算子、跨 DAG 依赖、SLA/任务丢失传感器以及自定义 Python 逻辑。
- 如果任务是短期的且吞吐量极高,首选无服务器的 Step Functions Express 或用于托管 ETL 操作的 Glue。
- 将繁重、长时间运行的任务保留在托管计算服务 (Glue/EMR/EKS) 中,仅将 MWAA 任务用作编排——避免在 MWAA 工作节点本身上运行大规模数据转换。
Airflow 中的错误处理:在 DAG 定义中使用任务重试和 retry_delay,设置 on_failure_callback 以进行通知或推送到 SQS 死信队列,并配置任务级别的 SLA 处理以触发修复 DAG。
用于无服务器编排的 AWS Step Functions
Step Functions 提供有状态的编排,使用基于 JSON 的 Amazon States Language,并与众多 AWS 服务广泛集成。在标准 (Standard) 和快速 (Express) 工作流之间进行选择:
- 标准工作流 (Standard Workflows):专为长期运行、持久的状态机(可持续数月至数年)而设计,具有精确一次 (唯一) 的执行语义、内置的执行历史记录以及每次执行的跟踪/日志记录。使用 aws stepfunctions start-execution –state-machine-arn arn:… –input ‘{“key”:“value”}’ 启动。
- 快速工作流 (Express Workflows):针对高吞吐量、低延迟、短时长的处理进行了优化,并且在大规模下具有成本效益;它们使用至少一次的执行语义,因此任务必须是幂等的或使用去重模式。
用例和决策标准:
- 当您需要可长期运行、需要仅一次执行语义的持久、可审计的工作流时,请使用标准工作流。
- 对于每秒数千次执行的事件驱动微编排,当持续时间短和成本效益很重要,并且您可以设计幂等任务或在下游进行去重时,请使用快速工作流。
错误处理和集成模式:
- 在 ASL 中使用 Retry 块,通过 ErrorEquals、IntervalSeconds、BackoffRate 和 MaxAttempts 来定义重试。
- 使用 Catch 块将失败重定向到备用分支或 Fail/Success 状态,并用错误详情填充 ResultPath 以进行诊断。
- 对于异步死信处理,将失败的消息推送到 SQS/SNS,或设计一个 Step Functions 模式,将错误负载发送到 SQS DLQ 进行离线处理。通过 LoggingConfiguration 和 TracingConfiguration 启用 CloudWatch Logs 和 X-Ray 跟踪以实现可观测性。
EventBridge Scheduler 和事件驱动的管道
EventBridge 为 cron 和一次性任务提供了丰富的事件路由和 Scheduler 功能。使用
undefined
创建基于计划的规则,并通过
undefined
附加目标。对于事件驱动(模式)路由,使用带有
undefined
的 put-rule 命令,将 S3 事件路由到 Lambda、Step Functions 或 SQS。
关键操作要点:
- EventBridge 支持计划表达式(cron 和 rate)。使用 rate 表达式时,请注意 EventBridge 规则有 5 分钟的最小间隔;如需更精细的粒度,请考虑使用 Step Functions 或轮询层。
- 使用 EventBridge Scheduler 进行一次性的、临时的未来调用和周期性计划;Scheduler 支持时区和灵活的按目标重试设置,并可以为无法投递的调用配置 SQS 死信队列。
- 对于高可靠性管道,附加像 Step Functions、Lambda 或 SQS 这样的目标,并配置按目标的重试策略和 DLQ。例如,
put-targets接受一个带有 SQS 队列 Arn 的DeadLetterConfig。
错误处理:配置特定于目标的重试次数和退避策略,对失败的投递使用 DLQ,并将 EventBridge 与 Step Functions 结合使用以处理复杂的错误和补偿事务。
常见陷阱与决策标准
- 错误:对非幂等任务使用 Express Workflows。正确方法:设计幂等性(去重键、幂等的 Lambda)或使用 Standard 工作流以实现“恰好一次”的语义。
- 错误:假设 Glue 条件触发器在失败时会触发。正确方法:显式创建 FAILED 状态的触发器,或在触发器的
Predicate中包含State=FAILED来路由错误。 - 错误:在公有子网中或在没有 NAT 的情况下部署 MWAA。正确方法:将 MWAA 放置在私有子网中,并为所需的服务访问提供 NAT 网关或 VPC 端点。
- 错误:期望 EventBridge 支持亚分钟级的计划。正确方法:记住 EventBridge 规则有 5 分钟的最小间隔;对于低于 5 分钟的需求,请使用 Step Functions 或 Lambda 计时器。
- 错误:缺乏跨服务的集中式重试/捕获策略。正确方法:标准化重试/退避策略(ASL 的
Retry、EventBridge 的重试配置、Airflow 的重试),并使用 DLQ 保存失败的事件以供手动/自动修复。 - 错误:用繁重的数据处理任务使 MWAA worker 超载。正确方法:在 MWAA 上只进行编排,在 Glue/EMR/EKS 上运行繁重的转换任务,并在任务之间传递指针(如 S3 路径)。
实践问题:Acme Retail 应对流量尖峰的每小时 ETL
Acme Retail 需要一个每小时运行的 ETL 流程,该流程运行 Glue 作业进行原始数据提取,一个带有 Python 算子的复杂数据浓缩 DAG,以及一个必须响应高频库存事件的短期 SKU 聚合。他们要求有稳健的重试机制和故障捕获能力。
- 使用 EventBridge 触发一个每小时的计划规则,该规则调用一个 Step Functions Standard 工作流来协调整个管道。
- 在 Step Functions 中,使用
Retry和Catch处理器来编排长时间运行的 Glue 作业(StartJobRun);在失败时,通过Catch块将任务路由到一个 SQS DLQ 和一个用于修复的 Lambda。 - 在 MWAA 中部署复杂的数据浓缩 DAG,并从 Step Functions 中通过 Airflow REST API 或通过向 SQS 放置 DAG 运行消息来调用它们;根据预期的并发量,通过
celery.worker_autoscale设置来调整 MWAA worker 的大小,并监控 CloudWatch 指标以进行调整。 - 对于高频库存事件,使用 EventBridge 事件模式规则将事件推送到一个带有幂等性密钥的 Express Step Function 或 Lambda,并使用一个由 SQS 支持的 DLQ 来吸收流量突发。
- 实施集中式监控(CloudWatch Logs/Metrics,用于 Step Functions 的 X-Ray),并针对 DLQ 积压增长和任务重试耗尽设置告警。
设计原理:此设计为每个需求都使用了正确的工具——Step Functions 用于持久化的跨服务编排和错误处理,MWAA 用于复杂的 DAG 逻辑,Glue 用于托管的 ETL,EventBridge 用于计划和反应式事件。它强制实施了幂等性和 DLQ,以构建符合 AWS 最佳实践的、有弹性且可观测的管道。
练习这些题目 → · 在 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.
通过考试 →