AWS Step Functions:如何等待触发器下所有Jobs完成后再继续
问题
我正在构建一个包含50多个Glue Jobs的AWS Step Functions状态机,部分Jobs存在依赖关系需要在工作流中体现,目前考虑了两种实现方案:
- 方案1:StartJob + 并行流:按依赖层级设置多组并行流,第一组执行无依赖的Jobs,第二组执行依赖第一组的任务,以此类推。通过勾选每个Job的「Wait for callback」选项,确保所有Job完成回调后才进入下一组并行流。
- 方案2:复用已有Glue触发器:现有Glue触发器已按依赖关系将Jobs分组为逻辑集合,想通过StartTrigger模块启动这些触发器,但遇到问题:为StartTrigger设置「Wait for callback」后,任务会在触发器启动后立即进入下一步,无法等待触发器下所有Jobs完成。
显然方案2更优,因为不想在状态机中手动创建60多个StartJob节点,但不知道如何实现。
尝试方案2时,编写的状态机代码出现「gluejob缺失元数据」错误,代码如下:
{ "Comment": "A description of my state machine", "StartAt": "Parallel", "States": { "Parallel": { "Type": "Parallel", "Branches": [ { "StartAt": "GetInitialTrigger", "States": { "GetInitialTrigger": { "Type": "Task", "Parameters": { "Name": "stepfunctionstest" }, "Resource": "arn:aws:states:::aws-sdk:glue:getTrigger", "Next": "Map", "OutputPath": "$.Trigger" }, "Map": { "Type": "Map", "ItemProcessor": { "ProcessorConfig": { "Mode": "INLINE" }, "StartAt": "Pass", "States": { "Pass": { "Type": "Pass", "Next": "Glue StartJobRun" }, "Glue StartJobRun": { "Type": "Task", "Resource": "arn:aws:states:::glue:startJobRun", "Parameters": { "JobName": "$.job_name", "Arguments.$":"$.Arguments" }, "End": true } } }, "InputPath": "$.Actions", "MaxConcurrency": 1, "End": true } } } ], "End": true } } }
解决方案
一、解决StartTrigger等待所有Jobs完成的问题
直接使用StartTrigger无法等待触发器关联的所有Jobs完成,可通过以下两种方式实现需求:
CloudWatch事件+Lambda回调
- 给触发器关联的所有Glue Jobs配置CloudWatch事件规则,当Job进入
SUCCEEDED或FAILED状态时触发Lambda函数。 - Lambda函数维护计数器,统计当前触发器下所有Job的完成状态,待全部Job完成后,调用Step Functions的
SendTaskSuccess/SendTaskFailure接口,通知状态机继续流程。 - 状态机中调用StartTrigger时开启「Wait for callback」,等待Lambda的回调信号。
- 给触发器关联的所有Glue Jobs配置CloudWatch事件规则,当Job进入
Lambda轮询Job状态
- 启动触发器后,添加一个Task节点调用Lambda函数,Lambda通过Glue的
GetJobRunsAPI轮询该触发器触发的所有Job运行状态。 - 当所有Job均完成(成功或失败)后,Lambda返回结果,状态机继续下一步。
- 启动触发器后,添加一个Task节点调用Lambda函数,Lambda通过Glue的
二、修复「gluejob缺失元数据」错误
你的代码存在两个关键问题:
- 参数路径不匹配:
GetInitialTrigger返回的Trigger.Actions中,Job名称的键是JobName(首字母大写),但代码中使用$.job_name(小写),导致无法正确获取Job名称。 - 冗余Pass节点:Pass节点无实际作用,可直接移除。
修复后的状态机代码如下:
{ "Comment": "State machine to run Glue jobs via trigger", "StartAt": "GetInitialTrigger", "States": { "GetInitialTrigger": { "Type": "Task", "Parameters": { "Name": "stepfunctionstest" }, "Resource": "arn:aws:states:::aws-sdk:glue:getTrigger", "Next": "Map", "OutputPath": "$.Trigger" }, "Map": { "Type": "Map", "ItemProcessor": { "ProcessorConfig": { "Mode": "INLINE" }, "StartAt": "Glue StartJobRun", "States": { "Glue StartJobRun": { "Type": "Task", "Resource": "arn:aws:states:::glue:startJobRun.sync", "Parameters": { "JobName.$": "$.JobName", "Arguments.$": "$.Arguments" }, "End": true } } }, "InputPath": "$.Actions", "MaxConcurrency": 1, "End": true } } }
注:使用arn:aws:states:::glue:startJobRun.sync作为资源,Step Functions会自动等待Glue Job完成,无需额外配置回调。
内容的提问来源于stack exchange,提问作者ben890
相关产品推荐
相关产品推荐

