详解在 Bedrock AgentCore 上用 SQS + Lambda + DynamoDB 构建事件驱动的环境感知代理,含人工介入工具设计。
大规模处理文档的团队对这套流程再熟悉不过:文件落入存储,有人注意到后打开每一个,决定需要做什么,再分发给相关人员审核。监控告警也是以同样的方式排队,等着相关人员去处理。手动分类所浪费的时间,正是环境智能体(ambient agent)要解决的操作痛点。想象一下,一份文档落入你的 Amazon S3 存储桶,几秒钟内你的 Jobs 页面就出现了一个作业,准备运行(或者如果你配置了的话,已经在运行了)。智能体分析文件,展示发现的结果,并在采取下一步操作前请求你批准。事件本身就成了提示词。这就是环境智能体:它响应事件流,当需要时通过单一的 ask_human 工具暂停等待人工输入,一旦人类回答便从中断处恢复。
图 1:AgentCore Runtime 上响应事件并暂停等待人工输入的环境智能体概述
如今大多数 AI 智能体体验遵循的是不同的模式:用户打开一个聊天界面,输入提示词,然后等待回复。这对一次性问题来说没问题,但它将智能体限制在一次一个对话中,而且需要有人先描述发生了什么,智能体才能对其采取行动。对于那些智能体应该对基础设施中发生的事件(文件上传、数据库变更、计划任务、系统告警)做出反应的场境来说,仅聊天模式就失效了。
环境智能体描述的是一种不同的范式,LangChain 等已经阐述过这种范式。环境智能体不再等待用户发起对话,而是监听事件流并对其采取行动,可能并行处理多个事件。它们并非 solely 由人类消息触发,多个智能体可以同时运行。最关键的是,它们并非完全自主:生产级设计需要仔细关注智能体何时暂停以与人类交互。当信号触发时,智能体执行其工作流,只有在需要澄清、批准或审查时才打断人类。这种人类参与式组件(human-in-the-loop)降低了将智能体部署到生产环境的风险,建立了用户信任,并让智能体能够通过反馈随时间学习和改进。
在 AWS 上运行的组织已经具备了事件驱动的基础设施:Amazon S3 事件通知、Amazon EventBridge 规则、AWS Lambda 触发器和 Amazon DynamoDB 流。缺失的一环是将这些事件源连接到能够推理发生了什么、采取行动,并在情况需要时引入人类的智能体。AWS Step Functions 这类全自动管道可以编排工作流,但无法对模糊性进行推理或提出澄清问题。基于聊天的智能体可以进行推理,但需要有人发起对话。环境智能体填补了这一空白。
Amazon Bedrock AgentCore 是一个用于大规模构建、连接和优化智能体的平台,支持任何框架或模型。AgentCore Runtime 提供了使这种模式成为可能的执行环境:支持长时间工作负载的基于容器的智能体托管、内置会话隔离,以及与 Amazon Bedrock 基础模型的集成。AgentCore Runtime 支持足够长的会话来覆盖此处所示的 signal → agent → human-in-the-loop(HITL)流程。参考实现将每次智能体轮次的上限设定为 Lambda 15 分钟超时,这在实践中绰绰有余。结合用于事件处理的 AWS Lambda 和用于状态管理的 Amazon DynamoDB,结果是一个完全无服务器的环境智能体平台。
在这篇文章中,我们将端到端地走过 Amazon Bedrock AgentCore 上的模式。你将理解:
ask_human 工具加上规范的响应信封如何足以支持全系列的人类参与式交互。在部署参考实现之前,确保你具备以下条件:
us-east-1(示例的默认配置针对该区域)。cdk bootstrap)。在我们深入架构之前,先了解一下是什么让环境智能体不同于典型的聊天机器人,以及文章的其余部分所依赖的构建模块会很有帮助:事件驱动触发模型、环境信号抽象,以及将它们绑在一起的单一人类参与式工具。
用户发起式智能体遵循请求-响应模式:
User → Prompt → Agent → Response → User
环境智能体遵循事件驱动模式:
Event → Signal → Agent → [Optional human interaction] → Action
关键区别在于触发机制。环境智能体由系统事件而非显式用户请求激活,这使它们天然适合文档处理管道、监控和告警、计划分析,以及需要沿路设置批准关卡的多步骤工作流。
环境信号是将事件源映射到智能体的配置。当事件发生时,平台自动为智能体创建一个作业。接下来发生什么取决于信号上的一个设置:
使用 autoExecute: false(默认值)时,作业以空闲状态出现在 Jobs 页面上,等待人工审查和运行。当信号可能在未知输入上触发,或者当智能体拥有高风险工具可用时,这是你想要的安全的先审查后执行流程。
使用 autoExecute: true 时,信号处理器直接将作业加入工作队列,智能体立即运行,只有在智能体本身调用 ask_human 时才会引入人工。这是完全自主的流程。
该模式涵盖几种信号事件源。参考示例附带了前两种。其余的是你通过编写新的处理程序 Lambda 函数和 Signals 页面上相应的表单字段来添加的扩展点:
jobType: "scheduled" 的作业驱动,而非 Signals 页面上的信号。环境智能体需要结构化的方式来与人类交互。在本示例中,智能体通过单一工具(ask_human)展示这些交互,并返回规范的响应信封。在该信封中,status 是 completed、interrupted 或 error 之一,匹配字段分别是 result、question 或 error。平台还通过每个响应线程传递 session_id 和 job_id,以便连续轮次可以关联。这些是关联元数据,不是你的智能体必须实现的核心契约的一部分。当智能体返回 interrupted 时,平台将作业移动到中断状态,并将其 requiresAction 标志设置为 true。参考 React 前端在 Jobs 页面的 Interrupted 选项卡上展示这些,并在每一行上显示警告指示器,因此没有单独的审查队列需要轮询。同一个 Jobs 视图展示待处理的问题、等待批准的建议操作、最终结果和失败的作业,让用户在一个地方看到他们的智能体正在做的一切,而无需监控多个聊天窗口或电子邮件线程。
同一套机制支撑多种提示模式,读者可以从更广泛的智能体文献中认出它们:Notify turn(通知轮次),智能体仅报告结果;Question turn(提问轮次),智能体请求澄清;Review turn(审核轮次),智能体提出行动并等待 APPROVE / REJECT / MODIFY;以及 Error turn(错误轮次),失败被捕获到作业记录中,由用户决定是否重试。这些都是智能体书写问题的约定,而非独立的运行时模式。在平台层面,只有一条代码路径和一个信封结构。
图 2:Notify、Question 和 Review 人工介入模式,均通过 ask_human 工具呈现
架构概述
平台由少量无服务器组件构成,通过事件管道串联在一起。本节先完整介绍端到端流程,再依次描述各组件。
事件在平台中端到端流转如下:Amazon S3 发出 s3:ObjectCreated 通知,Signal Processor Lambda 函数接收该通知。Signal Processor 在 ambient-signals 表上查询全局二级索引(GSI),找到与事件 bucket 匹配的任何信号,然后为每个匹配项创建一条作业记录。API 层(或调度器)将作业加入 Amazon SQS 队列。负责处理 API 路径的同一个 Job Execution Lambda 函数也通过附加的 SQS 事件源消费该队列,调用 Amazon Bedrock AgentCore Runtime 上的智能体,并将结果(及任何人机交互请求)写回 Amazon DynamoDB。从 Amazon S3 经 Amazon CloudFront 服务的 React 前端轮询一个小型的 Amazon API Gateway 和 Lambda 层以获取更新,并允许用户响应待处理的交互。
主要组件包括:
Amazon S3 配合事件通知作为信号的入口点,当文件上传时触发。Prefix 和 suffix 过滤器下推到 bucket 的通知配置中,使 Signal Processor 仅对可能匹配信号的事件进行调用。
Amazon SQS 将 API Gateway 请求与智能体调用解耦。作业执行队列持有待处理的工作。死信队列(DLQ)捕获工人在配置的重试次数后仍无法处理的消息。
三条管道 Lambda 函数将事件从摄入端传送到智能体(后面的管理平面由 API Gateway 背后的独立层级负责,下文详述):Signal Processor 将传入事件与已配置信号定义进行匹配并创建作业。Job Execution 在一个 Lambda 函数中有两条入口路径:一个是入队消息的 API 处理器,另一个是消费消息并使用作业上下文调用 AgentCore Runtime 的 SQS worker。Scheduler 以一分钟 cron 触发,将到期的定时作业加入同一 SQS 队列。
Signal Processor 将传入事件与已配置信号定义进行匹配并创建作业。
Job Execution 在一个 Lambda 函数中有两条入口路径:一个是入队消息的 API 处理器,另一个是消费消息并使用作业上下文调用 AgentCore Runtime 的 SQS worker。
Scheduler 以一分钟 cron 触发,将到期的定时作业加入同一 SQS 队列。
Amazon Bedrock AgentCore Runtime 在隔离容器中运行智能体代码,并支持长时间工作负载。
Amazon DynamoDB 存储智能体注册表、作业记录、环境信号定义、聊天线程、对话历史和 Powertools 幂等记录。对话消息通过 UpdateItem + list_append 原子追加,使并发写入者不会互相覆盖。
API Gateway 背后的五个管理平面 Lambda 函数暴露前端消费的 REST API(agent_management、job_management、signal_management、chat_management、conversation_management),另有一个 chat_execution worker Lambda 由 chat_management 异步调用,使聊天 API 调用立即返回。具体配置方式见生产部署一节。
从 Amazon S3 经 Amazon CloudFront 服务的 React 前端提供智能体管理 UI,用户在此监控作业、与智能体聊天、响应问题以及审核待处理操作。
图 3:事件流从 Amazon S3 经 Amazon SQS 和 AWS Lambda 到 AgentCore Runtime,状态存储在 Amazon DynamoDB,前端为 React
构建事件基础设施
在了解了架构之后,下一步是连接事件源,将 Amazon S3 上传转换为智能体作业。本节创建 bucket,将其事件通知指向 Signal Processor Lambda 函数,并介绍处理器如何将事件与已配置信号进行匹配。
设置 Amazon S3 信号触发器
创建 Amazon S3 bucket 并配置驱动 Signal Processor Lambda 函数的事件通知:
aws s3 mb s3://amzn-s3-demo-bucket-$(date +%s)
bucket 配置为直接将 s3:ObjectCreated:* 事件发送到 Signal Processor Lambda 函数。在本示例中,通知配置由 signal_management Lambda 函数在创建或更新信号时动态安装,因此为新 prefix 添加新信号无需重新部署。
Signal Processor Lambda 函数
Signal Processor 接收 Amazon S3 事件,在 DynamoDB 中查找匹配的信号定义,并为每个匹配项创建一个作业。其核心处理器形状如下面的示例所示。实际处理器在 backend/functions/multi_agent/signal_processor.py 中,还使用了 AWS Lambda Powertools 进行结构化日志记录和幂等,查询 signals 表上的 bucketName-signalId-index GSI,对每个信号上配置的 prefix 和 suffix 进行检查,并写入一条 signal_triggered 作业行。
def process_s3_signal(event):
"""Process S3 file-upload events and create agent jobs."""
for record in event["Records"]:
bucket = record["s3"]["bucket"]["name"]
key = record["s3"]["object"]["key"]
for signal in find_matching_signals(bucket, key):
create_agent_job(signal, {"bucket": bucket, "key": key})
信号配置数据模型
信号在 DynamoDB 中以下述结构存储:
{
"signalId": "sig-123abc",
"userId": "user-456def",
"agentId": "agent-789ghi",
"signalName": "Document Processor",
"signalType": "s3_file_upload",
"enabled": true,
"autoExecute": false,
"bucketName": "ambient-agent-documents",
"configuration": {
"bucketName": "ambient-agent-documents",
"prefix": "invoices/",
"suffix": ".pdf"
},
"triggerCount": 42,
"lastTriggered": "2026-04-15T10:30:00Z",
"createdAt": "2026-04-01T00:00:00Z"
}
通过 API 创建信号时,只需设置 configuration.bucketName。上例中显示的顶层 bucketName 由平台填充。DynamoDB GSI 分区键不能嵌套在 map 属性内,因此 signal_management 在每次写入时将 configuration.bucketName 镜像到顶层 bucketName 字段,使 bucketName-signalId-index GSI 能在每次 Amazon S3 事件时发散到匹配的信号。
autoExecute 标志是决定智能体是自主触发还是先由人工审核作业的唯一开关。智能体管理 UI 中的 Signals 表单在该设置旁将 autoExecute 暴露为一个复选框,与常见的 enabled 设置并列,因此更改行为只需快速编辑,无需触碰代码或数据库。
图 4:Signals 页面,信号定义上带有 autoExecute 开关
在 AgentCore Runtime 上部署智能体
Amazon Bedrock AgentCore Runtime 托管智能体容器,并通过 InvokeAgentRuntime API 暴露,worker Lambda 函数在每个作业时调用该 API。本节介绍示例附带的智能体布局、驱动它的配置,以及将 LangGraph 工具调用转换为平台响应信封的小型编排器。
智能体打包与结构
AgentCore Runtime 为智能体提供基于容器的执行环境,因此你可以引入任何 Python 智能体框架。示例采用模块化设计:
agent/
├── agent.py # Entry point
├── config.yaml # Agent configuration
├── requirements.txt # Dependencies
├── core/ # Platform integration
│ ├── agent_core.py # Main agent logic
│ ├── tool_factory.py # Config-driven tool creation
│ └── execution_control.py # Session state + loop detection
└── tools/ # Custom tools
├── calculator.py # Math operations
├── human_input.py # Human-in-the-loop
└── s3_reader.py # Exports list_s3_files + read_s3_file
config.yaml 定义了智能体行为、工具和系统提示词。本演示默认使用 Amazon Bedrock 上的 Anthropic Claude Sonnet 4.5,因为这种人机协同工作流依赖多步工具调用和长上下文推理能力,Claude Sonnet 4.5 非常适合。如果想切换到 Claude Haiku、Amazon Nova 或 Amazon Bedrock 上其他支持工具调用的模型(可用性因区域而异),只需修改配置文件中的 model_id 一行即可。max_iterations: 10 为图计算预留了约十次模型到工具的往返轮次空间,之后图执行会停止(编排层会将该值翻倍来计算 LangGraph 的递归限制,因为每次往返会穿越两个图节点),这足以覆盖示例中 S3、计算器和 ask_human 工具所需的多步工具调用。该文件还包含一个 execution: 代码块(循环检测器、熔断器、会话缓存阈值),此处为简洁起见已省略。完整文件请参见 agent/config.example.yaml。
# Amazon Bedrock Configuration
aws:
bedrock:
model_id: "us.anthropic.claude-sonnet-4-5-20250929-v1:0"
region_name: "us-east-1"
# Agent behavior
agent:
verbose: true
max_iterations: 10
handle_parsing_errors: true
# Tool configuration (enables config-driven tool composition)
tools:
calculator:
enabled: true
type: "calculator"
name: "calculator"
description: "Perform mathematical calculations."
human_input:
enabled: true
type: "human_input"
name: "ask_human"
description: "Request input or clarification from a human user."
s3_list:
enabled: true
type: "s3_list"
name: "list_s3_files"
description: "List files in an Amazon S3 bucket."
s3_reader:
enabled: true
type: "s3_reader"
name: "read_s3_file"
description: "Read and analyze files from Amazon S3."
prompts:
system_template: |
You are a helpful AI assistant with access to a set of tools.
Call tools only when they add value. Prefer concise answers
grounded in tool results.
核心智能体实现
智能体基于 langchain.agents.create_agent 构建,这是一个编译为 LangGraph 的工具调用智能体。选择 LangChain 作为编排层是因为它提供了预置的工具调用模式、广泛的开源生态系统集成以及许多团队已经熟悉的 API,而 AgentCore Runtime 则在底层提供托管服务、会话隔离和弹性扩缩容。两层是互补的。图计算外层的平台封装做了三件事:为当前轮次构建消息列表(包括对话历史,对于信号触发的工作流,还包括触发信号的 Amazon S3 存储桶和密钥,以便模型能够在不额外告知的情况下读取文件)、调用图计算,以及扫描工具输出中的 ask_human 哨兵值,将工具调用转换为中断响应。
最简形式下(完整版本在 agent/core/agent_core.py 中,还处理每会话历史缓存、循环检测和执行追踪),编排器的实现如下:
from bedrock_agentcore import BedrockAgentCoreApp
from langchain.agents import create_agent
from langchain_aws import ChatBedrock
from langchain_core.messages import ToolMessage
from core.tool_factory import create_tools_from_config, load_config
from tools.human_input import HUMAN_INPUT_SENTINEL
class Agent:
def __init__(self):
self.config = load_config()
self.llm = ChatBedrock(
model_id=self.config["aws"]["bedrock"]["model_id"],
region_name=self.config["aws"]["bedrock"]["region_name"],
)
# Tools are built from config.yaml so enabling a tool is a
# config change rather than a code change.
self.tools, _ = create_tools_from_config()
self.agent_graph = create_agent(
model=self.llm,
tools=self.tools,
system_prompt=self.config["prompts"]["system_template"],
)
def invoke(self, payload):
"""Run the agent graph for one turn. Simplified: the real method
also loads conversation history, binds execution state on a
ContextVar for tools to read, records an execution trace, and
threads ``session_id`` and ``job_id`` through on every response
so the platform can correlate continuation turns.
"""
session_id = payload.get("session_id", "default")
job_id = payload.get("job_id")
result = self.agent_graph.invoke(
{"messages": self._build_messages(payload)},
)
for msg in result["messages"]:
if isinstance(msg, ToolMessage):
content = str(msg.content)
if content.startswith(HUMAN_INPUT_SENTINEL):
question = content[len(HUMAN_INPUT_SENTINEL):].strip()
return {
"status": "interrupted",
"question": question,
"session_id": session_id,
"job_id": job_id,
}
return {
"status": "completed",
"result": self._extract_final_output(result["messages"]),
"session_id": session_id,
"job_id": job_id,
}
app = BedrockAgentCoreApp()
agent = Agent()
@app.entrypoint
def invoke(payload):
"""AgentCore Runtime entrypoint. Delegates to the Agent wrapper."""
return agent.invoke(payload)
使用提供的脚本将智能体部署到 AgentCore Runtime:
cd agent
# Configure AWS settings in .env and any overrides in config.yaml.
# Then build the container, push to Amazon ECR, and register it with
# Bedrock AgentCore Runtime.
./deploy_agent.sh
该脚本返回一个智能体运行时 Amazon Resource Name (ARN),后端将其存储在智能体注册表中,以便 Job Execution Lambda 函数可以调用它。
使用 DynamoDB 进行状态管理
Amazon DynamoDB 是所有需要超越单个 Lambda 调用生命周期的事物的系统记录:已注册的智能体、正在运行的工作流、赋予智能体跨轮次连续性的对话历史,以及 UI 读取的信号定义和聊天线程。下一小节将描述表布局、会话模型,以及 Job Execution Lambda 函数如何使用两者来驱动工作流完成。
状态层由 DynamoDB 支撑。本文涉及的核心表包括:
Agent 注册表