You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在AWS Step Functions中等待全部并行Lambda任务完成后进入下一状态?

设计并行Lambda执行流程(等待全部完成/回调后进入下一状态)

要实现你描述的流程,核心是利用AWS Step Functions的Parallel状态并行触发多个Lambda,同时结合Callback任务模式支持手动回调完成任务,确保所有任务结束后再进入state2。具体设计方案如下:

1. 用Parallel状态作为state1

Parallel状态的特性就是并行执行多个子分支任务,只有当所有分支都完成(无论是自动执行完毕还是通过回调标记完成),才会流转到下一个状态。你需要在这个状态里定义10个分支,每个分支对应一个Lambda调用。

2. 配置每个Lambda任务(支持自动完成或回调)

每个分支的核心是一个Task状态,配置如下:

  • 指定Lambda函数的ARN作为Resource
  • 如果需要支持回调,在Parameters里传入taskToken(通过内置变量$$.Task.Token获取),同时设置TimeoutSeconds(避免任务无限等待)
  • Lambda有两种完成方式:
    • 自动完成:Lambda执行逻辑结束后直接返回结果,Step Functions会自动标记该分支任务完成
    • 回调完成:Lambda不直接返回成功,而是保存收到的taskToken,之后在需要的时候调用SendTaskSuccess或SendTaskFailure API,手动通知Step Functions该任务完成

3. 状态机示例代码

下面是简化的状态机JSON结构,你可以根据实际需求补充完整10个Lambda分支:

{
  "Comment": "并行执行10个Lambda,全部完成后进入state2",
  "StartAt": "state1",
  "States": {
    "state1": {
      "Type": "Parallel",
      "Branches": [
        {
          "StartAt": "LambdaTask1",
          "States": {
            "LambdaTask1": {
              "Type": "Task",
              "Resource": "arn:aws:lambda:us-east-1:123456789012:function:MyLambda1",
              "Parameters": {
                "taskToken": "$$.Task.Token"
              },
              "TimeoutSeconds": 3600,
              "End": true
            }
          }
        },
        {
          "StartAt": "LambdaTask2",
          "States": {
            "LambdaTask2": {
              "Type": "Task",
              "Resource": "arn:aws:lambda:us-east-1:123456789012:function:MyLambda2",
              "Parameters": {
                "taskToken": "$$.Task.Token"
              },
              "TimeoutSeconds": 3600,
              "End": true
            }
          }
        }
        // 重复上述结构,添加LambdaTask3到LambdaTask10
      ],
      "Next": "state2"
    },
    "state2": {
      "Type": "Pass",
      "End": true
    }
  }
}

关键细节补充

  • 若选择回调模式,Lambda代码里需要处理taskToken参数,比如Python中可以这样调用回调API:
    import boto3
    import json
    
    stepfunctions = boto3.client('stepfunctions')
    
    def lambda_handler(event, context):
        task_token = event['taskToken']
        # 执行业务逻辑
        result = {"status": "success"}
        # 手动标记任务完成
        stepfunctions.send_task_success(
            taskToken=task_token,
            output=json.dumps(result)
        )
        # 注意:这里不需要返回结果给Step Functions,回调API会处理状态更新
    
  • Parallel状态会严格等待所有10个分支任务完成,无论每个任务是自动结束还是回调结束,都不会提前进入state2
  • 可以根据业务需求调整TimeoutSeconds,如果超过时限任务仍未完成,Step Functions会标记该任务失败,进而导致整个Parallel状态失败

内容的提问来源于stack exchange,提问作者zonyang

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 07:35:16