用 Amazon Bedrock AgentCore 构建事件驱动的 Ambient Agent:从 SQS 信号到人工审批

2026-10-02 26 预计阅读时间: 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.

预计阅读时间:12 分钟

传统智能体通常从一个聊天输入框开始:用户提问,模型回答。Ambient Agent 则反过来工作——它持续等待系统事件,例如 Amazon S3 新文件、定时任务或监控告警,并在没有聊天提示的情况下启动工作流。

Amazon Bedrock AgentCore 可以承载智能体运行逻辑,而 Amazon SQS、AWS Lambda 和 Amazon DynamoDB 则分别负责事件缓冲、任务执行和状态持久化。更关键的是,当智能体遇到高风险或信息不足的决策时,可以只暴露一个 ask_human 工具,把任务暂停到 Jobs 页面,等待人工确认后再继续。

不要把环境事件直接等同于模型提示词

一个可靠的 Ambient Agent 通常包含以下链路:

S3 / EventBridge / CloudWatch Alarm
                │
                ▼
              SQS
                │
                ▼
        Lambda Worker
                │
                ▼
      Bedrock AgentCore Runtime
                │
        ┌───────┴────────┐
        │                │
   自动完成任务      ask_human
                         │
                         ▼
                    DynamoDB
                         │
                         ▼
                    Jobs 页面
                         │
                    人工提交答案
                         │
                         ▼
                   SQS Resume 事件

SQS 在这里不只是连接服务的胶水。它将突发事件与智能体执行速度解耦,并提供重试、可见性超时和死信队列等保护。如果 S3 一次上传几千个对象,系统不必同时启动几千次长时间推理,而是由 Lambda 按设定的并发度消费。

进入队列的消息应该是稳定的业务事件,而不是直接拼接好的大段提示词。例如:

{
  "event_id": "evt-20250308-001",
  "event_type": "document.uploaded",
  "occurred_at": "2025-03-08T10:30:00Z",
  "resource": {
    "bucket": "incoming-documents",
    "key": "contracts/acme.pdf"
  },
  "policy": {
    "allow_external_action": false,
    "require_human_for": ["send_email", "approve_contract"]
  }
}

这种事件封装与具体智能体框架无关。Lambda 可以把它转换成 AgentCore 所需的调用格式,也可以在更换编排框架时继续沿用相同的队列协议。

ask_human 应该是暂停协议,而不只是一个问答工具

只给智能体一个人工协作工具,可以显著缩小控制面。智能体不需要知道 Slack、邮件或前端页面的实现,只需要在无法安全继续时调用:

{
  "name": "ask_human",
  "description": "Pause the current job and request a human decision.",
  "inputSchema": {
    "type": "object",
    "properties": {
      "question": { "type": "string" },
      "reason": { "type": "string" },
      "options": {
        "type": "array",
        "items": { "type": "string" }
      }
    },
    "required": ["question", "reason"]
  }
}

这里展示的是逻辑工具定义;实际注册方式需要按照所使用的 AgentCore 运行时和智能体框架调整。

工具被调用时,不应让 Lambda 一直阻塞等待。正确做法是:

  1. 将问题、上下文和恢复所需的信息写入 DynamoDB。
  2. 把任务状态改为 WAITING_HUMAN。
  3. 正常结束当前 SQS 消费。
  4. Jobs 页面展示待处理任务。
  5. 人工提交答案后,产生一个新的 human.answer_submitted 事件。
  6. Worker 使用保存的上下文恢复智能体。

推荐至少保存这些字段:

字段 用途
job_id 任务主键,也是幂等标识的一部分
status RUNNING、WAITING_HUMAN、COMPLETED 或 FAILED
source_event_id 追踪触发任务的原始事件
question 展示给审核人员的问题
reason 智能体为什么不能自动继续
options 可选答案或动作
checkpoint 恢复任务所需的会话 ID、步骤或业务上下文
answer 人工提交的结果
created_at / updated_at 排序、超时和审计

不要把完整提示词、机密文档或无裁剪的模型上下文全部塞入 DynamoDB。更稳妥的方式是保存 S3 引用、经过筛选的摘要以及恢复标识,并对敏感字段进行加密和访问控制。

可以直接改造的 DynamoDB 与 Python 骨架

下面示例假设已经配置 AWS CLI,并选择了目标 Region。它创建一张任务表,其中 status-created_at-index 可供 Jobs 页面查询待审核任务。

aws dynamodb create-table \
  --table-name AmbientAgentJobs \
  --attribute-definitions \
      AttributeName=job_id,AttributeType=S \
      AttributeName=status,AttributeType=S \
      AttributeName=created_at,AttributeType=S \
  --key-schema AttributeName=job_id,KeyType=HASH \
  --global-secondary-indexes '[{
    "IndexName":"status-created_at-index",
    "KeySchema":[
      {"AttributeName":"status","KeyType":"HASH"},
      {"AttributeName":"created_at","KeyType":"RANGE"}
    ],
    "Projection":{"ProjectionType":"ALL"},
    "ProvisionedThroughput":{"ReadCapacityUnits":5,"WriteCapacityUnits":5}
  }]' \
  --provisioned-throughput ReadCapacityUnits=5,WriteCapacityUnits=5

aws sqs create-queue --queue-name ambient-agent-resume

生产环境通常应使用按需计费、加密、备份和基础设施即代码;这里使用较短的 CLI 命令,是为了便于验证数据模型。

下面的 Python 模块实现了 ask_human、待办查询和人工回答提交。运行前安装 boto3,设置 TABLE_NAME 与 RESUME_QUEUE_URL。代码中的 checkpoint 是框架无关对象,可替换为 AgentCore 会话标识或自己的工作流状态。

import json
import os
from datetime import datetime, timezone

import boto3
from botocore.exceptions import ClientError

TABLE_NAME = os.environ.get('TABLE_NAME', 'AmbientAgentJobs')
RESUME_QUEUE_URL = os.environ['RESUME_QUEUE_URL']

dynamodb = boto3.resource('dynamodb')
table = dynamodb.Table(TABLE_NAME)
sqs = boto3.client('sqs')


def now_iso():
    return datetime.now(timezone.utc).isoformat()


def ask_human(job_id, source_event_id, question, reason,
              checkpoint, options=None):
    item = {
        'job_id': job_id,
        'status': 'WAITING_HUMAN',
        'source_event_id': source_event_id,
        'question': question,
        'reason': reason,
        'options': options or [],
        'checkpoint': checkpoint,
        'created_at': now_iso(),
        'updated_at': now_iso(),
    }

    try:
        table.put_item(
            Item=item,
            ConditionExpression='attribute_not_exists(job_id)'
        )
    except ClientError as exc:
        if exc.response['Error']['Code'] != 'ConditionalCheckFailedException':
            raise
        # SQS 可能重复投递。同一个 job_id 已存在时,将其视为幂等成功。

    return {
        'state': 'WAITING_HUMAN',
        'job_id': job_id,
        'message': 'The job is paused for human review.'
    }


def list_waiting_jobs(limit=50):
    response = table.query(
        IndexName='status-created_at-index',
        KeyConditionExpression='#status = :status',
        ExpressionAttributeNames={'#status': 'status'},
        ExpressionAttributeValues={':status': 'WAITING_HUMAN'},
        ScanIndexForward=True,
        Limit=limit,
    )
    return response.get('Items', [])


def submit_answer(job_id, answer, reviewer):
    updated = table.update_item(
        Key={'job_id': job_id},
        UpdateExpression=(
            'SET #status = :answered, answer = :answer, '
            'reviewer = :reviewer, updated_at = :updated'
        ),
        ConditionExpression='#status = :waiting',
        ExpressionAttributeNames={'#status': 'status'},
        ExpressionAttributeValues={
            ':waiting': 'WAITING_HUMAN',
            ':answered': 'ANSWERED',
            ':answer': answer,
            ':reviewer': reviewer,
            ':updated': now_iso(),
        },
        ReturnValues='ALL_NEW',
    )['Attributes']

    sqs.send_message(
        QueueUrl=RESUME_QUEUE_URL,
        MessageBody=json.dumps({
            'event_type': 'human.answer_submitted',
            'job_id': job_id,
            'answer': answer,
            'reviewer': reviewer,
            'checkpoint': updated['checkpoint'],
        }),
        MessageGroupId=job_id if RESUME_QUEUE_URL.endswith('.fifo') else None,
    ) if RESUME_QUEUE_URL.endswith('.fifo') else sqs.send_message(
        QueueUrl=RESUME_QUEUE_URL,
        MessageBody=json.dumps({
            'event_type': 'human.answer_submitted',
            'job_id': job_id,
            'answer': answer,
            'reviewer': reviewer,
            'checkpoint': updated['checkpoint'],
        }),
    )

    return updated

如果将其放入 Lambda,可以通过 API Gateway 暴露两个接口:

GET  /jobs?status=WAITING_HUMAN
POST /jobs/{job_id}/answer

Jobs 页面不必复杂。它只需要显示问题、原因、来源资源、允许的操作和审计信息,并调用回答接口。真正重要的是鉴权:审核人员必须通过 IAM Identity Center、Amazon Cognito 或组织现有身份系统登录,后端还应检查其是否有权处理对应任务。

示例代码在更新 DynamoDB 后再发送 SQS 消息,这两个操作并非原子事务。如果队列发送失败,任务会停留在 ANSWERED 状态但没有被恢复。生产系统可以使用 DynamoDB Streams 生成恢复事件,或者采用事务性 outbox 与定时修复任务,避免这个双写窗口。

恢复执行时,要防止智能体重复产生副作用

SQS 提供的是至少一次投递语义,因此同一个上传事件或人工回答可能被处理多次。仅靠“模型应该记得已经做过”并不可靠。Worker 应在调用 AgentCore 或外部系统前检查状态,并为副作用建立确定性的幂等键。

例如,可以使用:

idempotency_key = source_event_id + ':' + action_name + ':' + action_version

发送邮件、创建工单、更新数据库或调用支付接口时,都应记录这个键。重复消息到达时返回先前结果,而不是再次执行动作。

恢复事件也不应把人工回答直接当作新的自由文本命令。Worker 应加载原来的 checkpoint,把答案作为对应审核步骤的结构化结果,再继续原任务。这样可以降低提示注入和上下文错配风险。

上线前的检查清单

这类系统的挑战通常不在于让模型“自动运行”,而在于让自动运行可暂停、可恢复、可审计:

  • 为原始事件、任务和外部动作分别设计幂等键。
  • 为 SQS 配置死信队列,并监控消息年龄和失败次数。
  • 让 Lambda 可见性超时大于单次执行时间,长任务则拆分步骤。
  • 限制 AgentCore 运行角色、Lambda 和 Jobs API 的 IAM 权限。
  • 对 ask_human 设置超时、升级路径和任务所有者。
  • 在页面中展示证据与来源,不只展示模型结论。
  • 将高风险动作放在人工确认之后,而不是执行之后再通知。
  • 记录提示版本、工具版本、人工回答和最终副作用,满足审计需求。

Ambient Agent 的价值不是把聊天机器人藏到后台,而是把事件、模型推理、工具调用和人工判断组织成一条可靠的状态机。SQS 负责吸收不确定的事件流,DynamoDB 保存可恢复状态,Lambda 驱动执行,而 ask_human 则为自动化划出一条明确、安全的边界。


相关推荐