如何配置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类型,会自动生成任务令牌并传入指定的StoreTaskTokenLambda,此时流程会暂停,直到收到该令牌的回调请求。
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的读写权限。
StoreTaskTokenLambda需拥有DynamoDB的写入权限。- SNS回调Lambda需拥有DynamoDB的读取权限,以及调用Step Functions
SendTaskSuccess的权限。 - Step2的Lambda需拥有SNS消息发布权限。
内容的提问来源于stack exchange,提问作者Ajay
相关产品推荐
相关产品推荐

