在无服务器流水线中异步调用 Amazon Bedrock AgentCore Agent 的三种模式

2026-08-20 35 预计阅读时间: 1 分钟
来源: aws.amazon.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:11 分钟

当 Step Functions 调用 AI Agent 后需要等待较长时间,持续占用 Lambda 或其他计算资源会让流水线变贵,也会增加超时和重试处理的复杂度。针对 Amazon Bedrock AgentCore Agent,AWS 提供的思路是把“发起请求”和“等待结果”拆开,让 Agent 在后台处理,流水线只在真正需要时恢复执行。

本文围绕三种模式展开:任务令牌回调、直接服务集成和持久化函数。它们都能避免 Agent 处理期间的空转计算,但适用的控制粒度、集成复杂度和故障处理方式不同。

先明确异步调用的边界

典型流程可以拆成四个阶段:

  1. Step Functions 创建请求上下文,并生成或传递业务请求 ID。
  2. 流水线发起 AgentCore Agent 请求,然后进入等待状态。
  3. Agent 在独立的执行环境中处理请求,可能调用工具、访问数据或等待外部系统。
  4. Agent 完成后把结果写入约定位置,或者主动通知 Step Functions,流水线再继续执行。

关键点是:等待阶段不应依赖一个一直运行的 Lambda。状态机应该保存等待所需的信息,例如任务令牌、请求 ID、结果位置和超时截止时间。

模式一:任务令牌回调

任务令牌回调适合需要明确控制“何时恢复状态机”的场景。Step Functions 在等待状态中生成一个 task token,启动 Agent 的 Lambda 把 token 传给 Agent 或中间队列。Agent 完成后,回调方调用 SendTaskSuccessSendTaskFailure,状态机才会继续。

这个模式的优点是完成通知非常明确,适合长时间运行、结果大小有限或结果存储在 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:SendTaskSuccessstates: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

直接集成的难点在于接口契约。服务至少需要返回稳定的 requestIdstatus、错误码和结果位置,并支持幂等键。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 运行较久时继续保持低成本和可恢复性。


相关推荐