LangGraph + AWS:多 Agent 市场监测系统设计
展示用 LangGraph 和 AWS Bedrock 构建生产级多 agent 系统的完整架构。涉及编排、检查点恢复等核心技术。
展示用 LangGraph 和 AWS Bedrock 构建生产级多 agent 系统的完整架构。涉及编排、检查点恢复等核心技术。
随着人工智能应用从简单聊天机器人演进为复杂的自主系统,组织面临着新的挑战:需要编排复杂的多智能体工作流,以处理现实生产场景。传统的单智能体方法在处理需要专业知识、动态决策制定和强大错误恢复机制的复杂业务流程时往往力不从心。金融服务行业是这一挑战的典型例证。市场监测系统必须协调多个专业化的智能体来分析交易模式、调查可疑活动并生成全面的报告,同时要维持严格的合规性和可靠性标准。
解决方案结合了两个框架:用于宏观级工作流编排的 LangGraph 和用于智能体推理的 Strands。LangGraph 擅长管理多智能体协调的状态和有向图。它为您提供了对工作流执行和智能体间共享状态的细粒度控制。其中央持久层支持对生产至关重要的功能,包括人工参与式交互和基于检查点的强大故障恢复。同时,Strands Agent 充当工作流节点内的推理引擎。它提供模型无关的能力,可与多种大语言模型(LLM)提供商集成,同时保持灵活的工具集成和全面的可观测性。
随着亚马逊 Bedrock AgentCore 去年的发布,对于许多用例来说,AI 智能体解决方案的生产化可能会得到简化。这种组合为生产级 AI 智能体系统提供了坚实的基础,可以处理复杂的用例,同时帮助交付企业应用所需的基础设施可靠性和可观测性。
在这篇文章中,我们演示了如何在 AWS 基础设施上使用 LangGraph 和 Strands 来架构和部署多智能体 AI 系统。您将学到如何使用 LangGraph 的检查点系统实现状态驱动的工作流编排、集成 Strands 智能体用于专业化推理任务,以及使用 AgentCore 进行可扩展的生产部署。完整的解决方案已在 GitHub 上提供。
Strands Agent 运行在模型无关的架构上,可以适应您现有的基础设施而无需施加架构约束。该智能体实现了一个推理循环,持续评估工具输出并根据中间结果做出决策,这样您可以构建复杂的多步骤分析工作流。该框架包括全面的会话和状态管理,以及多个对话管理器,可防止您的上下文窗口溢出。
使用 Strands,您可以通过定义工具架构和访问模式来配置外部工具交互。对于我们的监测智能体,我们将数据发现与检索分离,以避免幻觉并加强对注入攻击的防御。我们使用工具如 get_report_list 和 get_report_schema 来查找报告,使用 run_report 来构建带有验证参数的 SQL 查询并运行它。
我们创建了带有以下工具和系统提示的 security_monitor 智能体:
from strands import Agent, tool
from strands.models.bedrock import BedrockModel
model = BedrockModel(
model_id="us.anthropic.claude-sonnet-4-6",
region_name="us-east-1",
max_tokens=16000,
additional_request_fields={
"thinking": {"type": "adaptive", "budget_tokens": 8000},
},
cache_prompt="default",
)
@tool
def get_report_list(agent_name: str) -> str:
"""Load the list of available reports for a specific agent.
Args:
agent_name: Name of the agent (e.g. 'security_monitor').
Returns:
str: JSON array of report objects with name and description.
"""
reports_data = load_agent_reports(agent_name)
return json.dumps(reports_data["reports"], indent=2)
@tool
def get_report_schema(report_name: str, query_intent: str) -> str:
"""Load column definitions for a report so you can build queries.
Args:
report_name: Name of the report (e.g. 'TradeActivity').
query_intent: Description of what data to extract.
Returns:
str: JSON object with parameters and column definitions.
"""
return json.dumps(load_json_report_definition(report_name), indent=2)
@tool
def run_report(
report_name: str,
filters: Dict[str, Any],
limit: Optional[int] = None,
) -> Dict[str, Any]:
"""Run a predefined report. The tool validates every filter against the
report's schema and builds a parameterised SQL query. The LLM never writes
raw SQL, so filter values cannot be injected into the query.
Args:
report_name: A report from `get_report_list` (e.g. 'TradeActivity').
filters: Equality filters keyed by column name, e.g.
{"symbol": "AAPL", "date": "2024-03-15"}.
limit: Optional row cap (1..10000).
Returns:
dict: {'success': bool, 'data': str (CSV), 'error': str or None}
"""
schema = load_json_report_definition(report_name)
allowed_columns = {c["name"] for c in schema["columns"]}
# Reject any filter field that is not in the report's allowed list of columns.
unknown = set(filters) - allowed_columns
if unknown:
raise ValueError(
f"Unknown filter field(s) {sorted(unknown)} for {report_name}. "
f"Allowed: {sorted(allowed_columns)}"
)
# Build SQL with named bind parameters.
where = " AND ".join(f"{field} = :{field}" for field in filters)
sql = f"SELECT * FROM {schema['reportName']}"
if where:
sql += f" WHERE {where}"
if limit is not None:
if not isinstance(limit, int) or not 1 <= limit <= 10_000:
raise ValueError("limit must be an integer in [1, 10000]")
sql += f" LIMIT {limit}"
return query_market_data(report_name=report_name, sql=sql, bind=filters)
security_monitor = Agent(
model=model,
system_prompt=SECURITY_MONITOR_PROMPT,
tools=[get_report_list, get_report_schema, run_report],
name="security_monitor",
)
LangGraph 通过三项核心能力为多智能体系统提供生产级编排,这些能力使其特别适合复杂的 AI 工作流。
基于图的状态机:LangGraph 将智能体工作流建模为有向图,其中节点代表包含智能体逻辑的函数,边确定执行流程。这种声明式方法将复杂的多步骤推理转化为可读、可维护的代码。图支持条件分支、并行执行和动态路由。这些能力对于工作流根据中间结果进行自适应的现实场景至关重要。
持久状态管理:该框架的检查点系统在每个节点执行后自动对完整的工作流状态进行快照。通过这些检查点,您可以从故障中平稳恢复并支持人工参与式交互。当分析师需要审查中间结果或发生错误时,系统从确切的检查点恢复。然后继续执行而不会丢失之前的工作。这种有状态的架构支持常见的模式,如多轮对话、迭代精化和跨越数小时或数天的长期调查。
生产可靠性:LangGraph 包含具有指数退避的内置重试策略,用于处理可能在您的系统中发生的限流或其他故障。它还通过 OpenTelemetry 提供全面的可观测性,方便与大多数可观测性应用集成。
让我们构建一个编排层,在我们的专业化 Strands 智能体之间路由查询。下面的代码定义了我们的工作流图:共享状态、选择调用哪个专业智能体的编排器、它们之间的条件路由,以及检查点支持的持久性,用于恢复和人工参与式审查。
from typing import TypedDict, Optional, List, Dict, Any
from langgraph.graph import END, StateGraph
from langgraph_checkpoint_aws import AgentCoreMemorySaver
class AgentState(TypedDict):
query_text: str
session_id: Optional[str]
agent_task_map: Optional[Dict[str, str]]
required_agents: Optional[List[str]]
current_agent_index: Optional[int]
# Each specialist writes its insights here
security_monitor_insights: Optional[Dict[str, Any]]
broker_monitor_insights: Optional[Dict[str, Any]]
risk_monitor_insights: Optional[Dict[str, Any]]
intel_analyst_insights: Optional[Dict[str, Any]]
synthesizer_insights: Optional[str]
SPECIALIST_NODES = {
"security_monitor": security_monitor_node,
"broker_monitor": broker_monitor_node,
"risk_monitor": risk_monitor_node,
"intel_analyst": intel_analyst_node,
}
def route_analysts(state: AgentState) -> str:
"""动态路由——按索引遍历 required_agents 列表。"""
required = state.get("required_agents", [])
index = state.get("current_agent_index", 0)
if not required:
return END
if index < len(required):
return required[index]
return "synthesizer"
# 构建图
workflow = StateGraph(AgentState)
workflow.add_node("orchestrator", orchestrator_node)
for name, node_fn in SPECIALIST_NODES.items():
workflow.add_node(name, node_fn)
workflow.add_node("synthesizer", synthesizer_node)
workflow.set_entry_point("orchestrator")
# 条件边——编排器和每个专家都通过 route_analysts 进行路由,
# 后者可以将任务移交给任意专家或综合器。
ALL_TARGETS = {name: name for name in SPECIALIST_NODES} | {
"synthesizer": "synthesizer",
END: END,
}
workflow.add_conditional_edges("orchestrator", route_analysts, ALL_TARGETS)
for name in SPECIALIST_NODES:
workflow.add_conditional_edges(name, route_analysts, ALL_TARGETS)
workflow.add_edge("synthesizer", END)
# AgentCoreMemorySaver 在每个节点执行后为状态创建检查点
checkpointer = AgentCoreMemorySaver(MEMORY_ID, region_name=REGION)
graph = workflow.compile(checkpointer=checkpointer)
市场分析 AI 智能体的 LangGraph 工作流编排
许多企业用例必须以严格、预定义的工作流为基础。完全依赖 LLM 的非确定性来执行正确的步骤会带来风险。然而,这些工作流中的某些具体步骤确实需要借助 LLM 的智能进行灵活推理。
LangGraph 与 Strands 的组合弥合了这一差距。你可以利用它构建这样的系统:由确定性编排承载局部的动态智能。
以下是这种架构组合解决复杂工作流挑战的方式:
“节点”级智能: LangGraph 定义工作流的高层编排。Strands AI 智能体被放置在特定节点中,用于处理复杂工作流中可能需要 LLM 分析或应对歧义的环节。它们只在确实需要这种灵活性的地方应用自主推理和工具调用。
通过节点 AI 智能体实现上下文隔离: 单体 AI 智能体很容易忘记自身的指令。通过将各个 Strands AI 智能体放入相互独立的 LangGraph 节点,可以将内存划分为不同的隔离区。每个 AI 智能体管理自身高度聚焦的上下文和工具历史记录,而 LangGraph 则维护一个结构化的全局会话状态,不同的 AI 智能体可以独立更新和引用该状态。
强化编排能力: 从本质上讲,LangGraph 是一种底层路由和编排工具。通过嵌入 Strands,你可以立即将一个全面的企业级 AI 智能体框架接入图中。这样既能获得 LangGraph 强大的路由能力,也能获得 Strands 原生的 Model Context Protocol(MCP)集成、引导控制、安全护栏和评估能力。
具体而言,每个专家都是一个 LangGraph 节点。节点会启动一个全新的 Strands AI 智能体,该 AI 智能体拥有自己的系统提示词、工具和隔离上下文。它执行编排器分配的任务并返回状态更新。LangGraph 将这部分更新合并到其他节点读取的共享状态中。以下是安全监控节点:
async def security_monitor_node(state: AgentState) -> AgentState:
"""
Single-day activity analyst agent that assesses price, volume and tick-level trades
"""
agent = Agent(
name="security_monitor",
model=analyst_model,
system_prompt=SECURITY_MONITOR_PROMPT,
tools=[get_report_list, get_report_schema, run_report],
callback_handler=None,
)
# Pull this node's task from the shared state the orchestrator populated.
task = state.get("agent_task_map", {}).get("security_monitor", state["query_text"])
# Strands runs its own reasoning + tool loop; collect the agent's final text.
chunks = []
async for event in agent.stream_async(task):
if "data" in event:
chunks.append(event["data"])
result = "".join(chunks)
# Return shared state updates
return {
"security_monitor_insights": {"task": task, "business_insights": result},
"current_agent_index": state.get("current_agent_index", 0) + 1,
}
Amazon Bedrock AgentCore 提供完全托管的服务,用于大规模部署和运行 AI 智能体,在提供生产级能力的同时减轻基础设施管理负担。
运行时部署: AgentCore runtime 是 Amazon Bedrock AgentCore 的一项能力,可以通过极少的配置将本地 AI 智能体代码转换为云原生部署。该服务与框架无关,可以直接配合 LangGraph 和 Strands 使用。它为动态 AI 智能体工作负载提供专用基础设施,包括支持长时间调查任务的扩展运行时、支持交互式工作流的低延迟执行,以及根据需求自动扩缩容的能力。
使用 AgentCore Python SDK(Amazon Bedrock AgentCore 的一项能力)和 starter toolkit 部署 LangGraph 编排器与 Strands AI 智能体。运行时会自动处理容器化、网络和计算资源预置。AgentCore 承担容器编排、扩缩容和会话管理等无差异化的繁重工作。
# api.py --- AgentCore runtime entry point
from bedrock_agentcore.runtime import BedrockAgentCoreApp
from src.agents import Workflow
app = BedrockAgentCoreApp()
workflow = Workflow()
@app.entrypoint
async def market_surveillance_workflow(payload):
"""Invoked by AgentCore for each request. Yields streaming chunks."""
prompt = payload.get("prompt")
session_id = payload.get("session_id", "default-session")
actor_id = payload.get("actor_id", "default-actor")
async for chunk in workflow.stream_query(
session_id=session_id, prompt=prompt, actor_id=actor_id
):
yield chunk
if __name__ == "__main__":
app.run()
# Deploy with the AgentCore starter toolkit
from bedrock_agentcore_starter_toolkit import Runtime
runtime = Runtime()
runtime.configure(
entrypoint="api.py",
auto_create_execution_role=True,
auto_create_ecr=True,
requirements_file="requirements.txt",
region="us-east-1",
agent_name="market_surveillance_workflow",
)
result = runtime.launch()
print(f"Agent ARN: {result.agent_arn}")
# Invoke the deployed agent
import boto3, json
client = boto3.client("bedrock-agentcore", region_name="us-east-1")
response = client.invoke_agent_runtime(
agentRuntimeArn=result.agent_arn,
qualifier="DEFAULT",
payload=json.dumps({
"prompt": "What caused the AAPL price spike at 11:00 AM on March 15, 2024?",
"session_id": "session-001",
"actor_id": "analyst-jane",
}),
)
# AgentCore returns a server-sent-events stream. Parse it:
for raw in response["response"].iter_lines():
if not raw:
continue
line = raw.decode("utf-8") if isinstance(raw, bytes) else raw
if not line.startswith("data: "):
continue
try:
chunk = json.loads(line[6:])
if isinstance(chunk, str) and chunk.startswith("data: "):
chunk = json.loads(chunk[6:])
except json.JSONDecodeError:
continue # malformed chunk --- skip, don't crash
if isinstance(chunk, dict) and chunk.get("type") == "text":
print(chunk["content"], end="")
内存集成: LangGraph 只需几行代码,即可通过 langgraph-checkpoint-aws 包与 AgentCore memory(Amazon Bedrock AgentCore 的一项能力)集成,同时提供短期检查点持久化和智能长期记忆检索功能。
本示例使用的 AgentCoreMemorySaver 类负责处理包含用户消息、AI 响应、图执行状态和元数据的检查点对象。每个节点执行完毕后,LangGraph 都会自动将检查点保存到 AgentCore memory。无需管理 Amazon DynamoDB 表,也无需实现自定义序列化逻辑,即可获得有状态对话和工作流恢复能力。
AgentCoreMemoryStore 类提供智能记忆能力,AgentCore 会自动从对话中提取洞察、摘要和用户偏好。AI 智能体可以在未来的交互中搜索这些记忆,从而提供随时间推移不断改进的个性化体验。这解决了 AI 智能体无状态这一根本性挑战:每次交互都建立在先前知识的基础上,而不是从头开始。
import boto3, time
REGION = "us-east-1"
control_client = boto3.client("bedrock-agentcore-control", region_name=REGION)
response = control_client.create_memory( name="MarketSurveillanceMemory", description="Memory for market surveillance multi-agent workflow.", eventExpiryDuration=90, # days ) MEMORY_ID = response["memory"]["id"] print(f"Memory ID: {MEMORY_ID}")
deadline = time.time() + 600 while True: status = control_client.get_memory(memoryId=MEMORY_ID)["memory"]["status"] if status == "ACTIVE": break if status == "FAILED" or time.time() >= deadline: raise RuntimeError(f"Memory {MEMORY_ID} is {status!r} (expected ACTIVE)") time.sleep(10)
```python
checkpointer = AgentCoreMemorySaver(MEMORY_ID, region_name=REGION)
graph = workflow.compile(checkpointer=checkpointer)
调用时,将 thread_id 和 actor_id 传递给智能体。它们分别是用户和会话的唯一标识符:
config = {
"configurable": {
"thread_id": "surveillance-session-001",
"actor_id": "analyst-jane",
}
}
response = await graph.ainvoke(
{"query_text": "Which brokers were most active?"},
config=config,
)
可观测性与运维:AgentCore 通过与 Amazon CloudWatch 和 AWS X-Ray 集成,提供内置的可观测性功能,用于捕获智能体执行追踪、工具调用和性能指标。该服务提供用于监控智能体行为、识别瓶颈和优化成本的仪表板。结合 LangGraph 的 OpenTelemetry 事件,你可以获得全面的可见性,涵盖从高层工作流编排到单次 LLM 调用和推理步骤的各个层面。
本文介绍了如何将 LangGraph 强大的工作流编排能力与 Strands 的智能化智能体推理能力相结合,并将其部署在 AWS 基础设施上,从而构建可用于生产环境的多智能体 AI 系统。
我们探讨的混合架构展示了 LangGraph 如何擅长宏观层面的编排(管理智能体协调、状态持久化和工作流恢复),而 Strands 则在各个节点内部提供细粒度的推理引擎。通过这种关注点分离,你可以构建能够处理复杂业务流程的精密系统,同时借助基于检查点的恢复机制和全面的可观测性来确保生产环境的可靠性。
你可以考虑将这一架构扩展到其他复杂的编排场景,例如文档处理流水线、客户服务自动化或合规监控系统。Strands 与模型无关的特性结合 LangGraph 的状态管理能力,使这一模式对于同时要求灵活性和可靠性的企业应用尤其有价值。
若要深入了解技术实现细节并自行构建市场分析智能体,请参阅 GitHub 仓库。