当 Step Functions 调用 AI Agent 后需要等待较长时间,持续占用 Lambda 或其他计算资源会让流水线变贵,也会增加超时和重试处理的复杂度。针对 Amazon Bedrock AgentCore Agent,AWS 提供的思路是把“发起请求”和“等待结果”拆开,让 Agent 在后台处理,流水线只在真正需要时恢复执行。
本文围绕三种模式展开:任务令牌回调、直接服务集成和持久化函数。它们都能避免 Agent 处理期间的空转计算,但适用的控制粒度、集成复杂度和故障处理方式不同。
先明确异步调用的边界
典型流程可以拆成四个阶段:
- Step Functions 创建请求上下文,并生成或传递业务请求 ID。
- 流水线发起 AgentCore Agent 请求,然后进入等待状态。
- Agent 在独立的执行环境中处理请求,可能调用工具、访问数据或等待外部系统。
- Agent 完成后把结果写入约定位置,或者主动通知 Step Functions,流水线再继续执行。
关键点是:等待阶段不应依赖一个一直运行的 Lambda。状态机应该保存等待所需的信息,例如任务令牌、请求 ID、结果位置和超时截止时间。
模式一:任务令牌回调
任务令牌回调适合需要明确控制“何时恢复状态机”的场景。Step Functions 在等待状态中生成一个 task token,启动 Agent 的 Lambda 把 token 传给 Agent 或中间队列。Agent 完成后,回调方调用 SendTaskSuccess 或 SendTaskFailure,状态机才会继续。
这个模式的优点是完成通知非常明确,适合长时间运行、结果大小有限或结果存储在 S3/DynamoDB 中的任务。代价是需要管理 token 的安全传递、幂等回调和超时清理。
下面是一个可以改造的 Amazon States Language 示例。示例中的 AgentCore 调用使用占位 Lambda 表示,实际项目中可以把它替换成负责调用 AgentCore Runtime 的函数。回调 Lambda 只负责把最终结果提交给 Step Functions。
Comment: Asynchronous AgentCore invocation with a task token
StartAt: StartAgent
States:
StartAgent:
Type: Task
Resource: arn:aws:states:::lambda:invoke.waitForTaskToken
TimeoutSeconds: 3600
Parameters:
FunctionName: arn:aws:lambda:us-east-1:123456789012:function:start-agentcore-request
Payload:
taskToken.$: $$.Task.Token
executionArn.$: $$.Execution.Id
input.$: $
ResultPath: $.agentResult
Next: ContinuePipeline
Catch:
- ErrorEquals: [States.Timeout]
Next: AgentTimedOut
ContinuePipeline:
Type: Pass
Parameters:
status: completed
agentResult.$: $.agentResult
End: true
AgentTimedOut:
Type: Fail
Error: AgentCoreTimeout
Cause: Agent did not complete before the workflow timeout
启动函数需要把 token 和请求信息交给异步处理器。下面的 Python 代码演示了一个最小实现,其中 submit_to_agentcore 是项目需要替换的适配层。它可以把请求放入 SQS,也可以调用负责 AgentCore 的内部服务。
import json
import os
import boto3
sqs = boto3.client("sqs")
QUEUE_URL = os.environ["AGENT_QUEUE_URL"]
def submit_to_agentcore(message: dict) -> None:
# Replace this adapter with the AgentCore Runtime invocation used by your account.
sqs.send_message(
QueueUrl=QUEUE_URL,
MessageBody=json.dumps(message),
)
def lambda_handler(event, context):
task_token = event["taskToken"]
request = {
"executionArn": event["executionArn"],
"taskToken": task_token,
"input": event.get("input", {}),
}
submit_to_agentcore(request)
return {"accepted": True}
完成处理的消费者需要确保回调只执行一次。实际实现通常会把 taskToken 和请求状态写入 DynamoDB,使用条件更新防止重复成功或重复失败,然后调用:
import json
import boto3
sfn = boto3.client("stepfunctions")
def report_success(task_token: str, result: dict) -> None:
sfn.send_task_success(
taskToken=task_token,
output=json.dumps(result),
)
不要让客户端直接持有 task token。它相当于恢复状态机所需的敏感凭证,应只在受信任的后端、队列和 Agent 回调服务之间传递,并通过 IAM 限制 states:SendTaskSuccess 和 states:SendTaskFailure 的权限范围。
模式二:直接服务集成
直接服务集成适合 AgentCore 已经暴露出可供 Step Functions 调用的 AWS API、HTTP 接口或内部服务接口,并且团队希望减少 Lambda 胶水代码的情况。Step Functions 直接完成请求提交、参数映射和错误路由,服务本身负责异步处理。
可以这样实践:把状态机拆成“提交请求”和“查询结果”两个状态。提交状态只拿到一个 requestId,随后使用 Wait 配合查询状态,直到结果完成或达到最大轮询次数。这个方案不会在 Wait 期间占用 Lambda,但轮询次数较多时会增加状态转换成本,因此应设置合理的等待间隔,并优先使用事件通知替代高频轮询。
下面的状态机是接口契约示例,arn:aws:states:::apigateway:invoke 和 URL 参数需要根据实际 AgentCore 接入方式调整:
StartAt: SubmitRequest
States:
SubmitRequest:
Type: Task
Resource: arn:aws:states:::apigateway:invoke
Parameters:
ApiEndpoint: agent-adapter.example.internal
Method: POST
Stage: prod
Path: /agent-runs
RequestBody:
prompt.$: $.prompt
correlationId.$: $$.Execution.Id
ResultPath: $.submission
Next: WaitBeforeQuery
WaitBeforeQuery:
Type: Wait
Seconds: 15
Next: QueryResult
QueryResult:
Type: Task
Resource: arn:aws:states:::apigateway:invoke
Parameters:
ApiEndpoint: agent-adapter.example.internal
Method: GET
Stage: prod
Path.$: States.Format('/agent-runs/{}', $.submission.requestId)
ResultPath: $.query
Next: IsComplete
IsComplete:
Type: Choice
Choices:
- Variable: $.query.status
StringEquals: completed
Next: ContinuePipeline
- Variable: $.query.status
StringEquals: failed
Next: AgentFailed
Default: WaitBeforeQuery
ContinuePipeline:
Type: Pass
End: true
AgentFailed:
Type: Fail
Error: AgentCoreRequestFailed
直接集成的难点在于接口契约。服务至少需要返回稳定的 requestId、status、错误码和结果位置,并支持幂等键。Step Functions 可能因为网络错误重试提交请求,因此请求头或请求体中的 correlationId 不能只用于日志,也应该参与服务端去重。
模式三:持久化函数
持久化函数适合需要编排多个长时间步骤、暂停后等待外部事件,或者希望把重试和恢复逻辑封装在函数工作流中的场景。它的核心不是让一个函数持续占用计算资源,而是让函数的执行历史、等待点和恢复点被持久化。
在这种模式下,持久化函数可以负责以下工作:
- 创建 AgentCore 请求并保存业务上下文。
- 等待 Agent 完成事件、人工审批或外部系统通知。
- 在恢复后读取 Agent 结果,并决定是否重试、补偿或继续下一步。
- 为每个阶段设置独立的超时和错误策略。
它比单个 Lambda 中的同步等待更适合复杂流程,但也引入了新的运行时和调试模型。需要确认所选持久化函数能力是否覆盖当前区域、部署方式、日志关联、最大执行时长以及幂等语义。不要仅因为函数名称包含“durable”就假设所有外部 SDK 调用都会自动变成可恢复操作。
选择建议与生产检查清单
可以按下面的决策方式选择:
- 需要 Agent 完成后精准唤醒同一个 Step Functions 状态:优先任务令牌回调。
- 已有稳定的异步 Agent 服务接口,希望减少编排代码:考虑直接服务集成。
- 一个任务包含多个等待点、人工介入和补偿步骤:考虑持久化函数。
上线前至少检查这些项目:
- 为每个请求定义
requestId和幂等键,避免状态机重试产生重复 Agent 任务。 - 为提交、等待、回调和最终结果分别配置超时,而不是只设置一个总超时。
- 结果较大时,把结果写入 S3 或 DynamoDB,状态机只传递引用。
- 对 task token、Agent 凭证和结果数据使用最小 IAM 权限与加密传输。
- 为回调失败、Agent 失败、状态机超时和重复回调建立可观测性指标。
- 明确哪些错误可以重试,哪些错误需要人工处理或补偿。
异步化并不只是把同步 API 改成后台任务。真正的收益来自清晰的状态边界:提交请求时记录什么,等待期间由谁负责通知,完成后如何验证结果,以及重复消息和超时如何收敛。把这几个问题设计清楚,Step Functions 才能在 Agent 运行较久时继续保持低成本和可恢复性。