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

如何配置AWS Step Functions,让Step3等待SNS消息触发?

实现AWS Step Functions中Step3等待SNS消息的方案

要让Step3仅在收到SNS消息时触发,核心是利用Step Functions的任务令牌(Task Token)回调机制——让流程在Step2执行完成后进入暂停状态,直到SNS消息触发的Lambda调用SendTaskSuccess API唤醒流程,继续执行Step3。

以下是具体实现步骤:

1. 编写Step Functions流程定义

流程中需要添加一个等待回调的任务节点,用于生成任务令牌并暂停流程。示例JSON如下:

{
  "StartAt": "Step1",
  "States": {
    "Step1": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:REGION:ACCOUNT_ID:function:Step1Lambda",
      "Next": "Step2"
    },
    "Step2": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:REGION:ACCOUNT_ID:function:Step2Lambda",
      "Parameters": {
        "executionId.$": "$$.Execution.Id"
      },
      "Next": "WaitForSNSNotification"
    },
    "WaitForSNSNotification": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke.waitForTaskToken",
      "Parameters": {
        "FunctionName": "arn:aws:lambda:REGION:ACCOUNT_ID:function:StoreTaskToken",
        "Payload": {
          "taskToken.$": "$$.Task.Token",
          "executionId.$": "$$.Execution.Id"
        }
      },
      "Next": "Step3"
    },
    "Step3": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:REGION:ACCOUNT_ID:function:Step3Lambda",
      "End": true
    }
  }
}

说明:WaitForSNSNotification节点使用lambda:invoke.waitForTaskToken类型,会自动生成任务令牌并传入指定的StoreTaskToken Lambda,此时流程会暂停,直到收到该令牌的回调请求。

2. 实现StoreTaskToken Lambda

这个Lambda的作用是将任务令牌与流程的Execution ID关联存储到DynamoDB,方便后续SNS触发的Lambda查询。示例Python代码:

import boto3

dynamodb = boto3.resource('dynamodb')
table = dynamodb.Table('TaskTokenStore')

def lambda_handler(event, context):
    # 存储令牌和Execution ID到DynamoDB
    table.put_item(
        Item={
            'executionId': event['executionId'],
            'taskToken': event['taskToken']
        }
    )
    return {}

3. 修改Step2的Lambda逻辑

Step2的Lambda在完成递归调用后,发送SNS消息时要携带流程的Execution ID,方便后续回调Lambda定位任务令牌。示例Python代码:

import boto3
import json

sns_client = boto3.client('sns')
TARGET_TOPIC_ARN = 'arn:aws:sns:REGION:ACCOUNT_ID:StepCompletionTopic'

def lambda_handler(event, context):
    execution_id = event['executionId']
    
    # 执行递归调用逻辑,直到任务完成
    # ... 你的递归代码 ...
    
    # 发送SNS消息,附带Execution ID
    sns_client.publish(
        TopicArn=TARGET_TOPIC_ARN,
        Message=json.dumps({'executionId': execution_id, 'status': 'completed'})
    )
    
    return {'status': 'sns_sent'}

4. 实现SNS触发的回调Lambda

创建一个Lambda并订阅Step2使用的SNS主题,该Lambda负责从DynamoDB取出任务令牌,调用Step Functions的SendTaskSuccess API唤醒流程。示例Python代码:

import boto3
import json

dynamodb = boto3.resource('dynamodb')
table = dynamodb.Table('TaskTokenStore')
sf_client = boto3.client('stepfunctions')

def lambda_handler(event, context):
    # 解析SNS消息中的Execution ID
    sns_message = json.loads(event['Records'][0]['Sns']['Message'])
    execution_id = sns_message['executionId']
    
    # 从DynamoDB获取对应的任务令牌
    db_response = table.get_item(Key={'executionId': execution_id})
    task_token = db_response['Item']['taskToken']
    
    # 调用SendTaskSuccess唤醒流程
    sf_client.send_task_success(
        taskToken=task_token,
        output=json.dumps({'message': 'SNS notification received'})
    )
    
    # 清理DynamoDB中的令牌(可选)
    table.delete_item(Key={'executionId': execution_id})
    
    return {'status': 'flow_resumed'}

5. 配置必要的权限

确保所有组件的权限配置正确:

  • Step Functions执行角色需拥有调用所有涉及Lambda的权限,以及DynamoDB的读写权限。
  • StoreTaskToken Lambda需拥有DynamoDB的写入权限。
  • SNS回调Lambda需拥有DynamoDB的读取权限,以及调用Step Functions SendTaskSuccess的权限。
  • Step2的Lambda需拥有SNS消息发布权限。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 09:45:32