AWS Machine Learning Blog

Market surveillance agent with LangGraph and Strands on AgentCore

8.5内容质量
Market surveillance agent with LangGraph and Strands on AgentCore

TL;DR · AI 摘要

AWS博客展示如何用LangGraph和Strands框架构建金融监管代理系统,结合AgentCore实现多代理工作流编排与智能推理。

核心要点

  • LangGraph通过检查点系统实现多代理工作流状态管理
  • Strands框架支持模型无关的智能代理推理架构
  • AgentCore提供可扩展的生产级代理系统部署方案

结构提纲

按章节快速跳转。

  1. 阐述金融监管场景下传统单代理系统的技术局限性

  2. 对比LangGraphStrands在工作流编排与智能推理层面的互补性

  3. AgentCore集成方案

    说明Amazon Bedrock AgentCore如何简化生产部署流程

  4. 展示基于AWS基础设施的多代理系统部署拓扑结构

  5. 解析Strands的模型无关架构与LangGraph的检查点机制

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • 金融监管代理系统架构
    • LangGraph
      • 工作流编排
      • 状态管理
    • Strands
      • 模型无关推理
      • 工具集成
    • AgentCore
      • 生产部署
      • 可观测性

金句 / Highlights

值得收藏与分享的关键句。

#LangGraph#Strands#AgentCore#多代理系统#金融监控
打开原文

基于 LangGraph 和 Strands 的市场监控代理 | 人工智能

基于 LangGraph 和 Strands 的市场监控代理

随着人工智能应用从简单的聊天机器人演进为复杂的自主系统,组织在协调能够处理现实生产场景的复杂多代理工作流程时面临新的挑战。传统单代理方法在处理需要专业知识、动态决策和强大错误恢复机制的复杂业务流程时往往力不从心。金融服务行业就是这种挑战的典型代表。市场监控系统必须协调多个专业代理来分析交易模式、调查可疑活动并生成全面报告,同时保持严格的合规性和可靠性标准。

该解决方案结合了两种框架:LangGraph 用于宏观层面的工作流程编排,Strands 用于智能代理推理。LangGraph 擅长管理状态和有向图以实现多代理协调。它为您提供对工作流程执行和代理间共享状态的细粒度控制。其核心持久化层支持生产环境中至关重要的功能,包括人机交互和基于检查点的稳健故障恢复机制。同时,Strands Agent 作为单个工作流程节点中的推理引擎。它提供与各种大语言模型(LLM)供应商集成的模型无关能力,同时保持灵活的工具集成和全面的可观测性。

随着 Amazon Bedrock AgentCore 在去年的发布,许多用例的代理解决方案生产化可能变得更加简单。这种组合为生产就绪的代理人工智能系统提供了坚实的基础,能够处理复杂用例,同时满足企业应用程序对基础设施可靠性和可观测性的需求。

在本文中,我们将演示如何使用 LangGraph 和 Strands 在 AWS 基础设施上构建和部署多代理人工智能系统。您将学习如何使用 LangGraph 的检查点系统实现状态驱动的工作流程编排,如何集成 Strands 代理执行专业推理任务,以及如何使用 AgentCore 进行可扩展的生产部署。完整的解决方案可在 GitHub 上获取。

Strands:智能代理推理

Strands Agent 基于模型无关架构构建,能够适应现有基础设施而不会强加架构限制。该代理实现了代理推理循环,持续评估工具输出并基于中间结果做出决策,因此您可以构建复杂的多步骤分析工作流程。该框架包含全面的会话和状态管理以及多个对话管理器,可防止上下文窗口溢出。

使用 Strands,您可以通过定义工具模式和访问模式来配置外部工具交互。对于我们的监控代理,我们将数据发现与数据检索分离,以避免幻觉并增强对注入攻击的防护。我们使用 get_report_list 和 get_report_schema 工具查找报告,并使用 run_report 工具构建带有验证参数的 SQL 查询并执行。

我们创建了一个带有以下工具和系统提示的安全监控代理:

code
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:
    """加载特定代理可用报告的列表。

    参数:
        agent_name: 代理名称(例如 'security_monitor')。

    返回:
        str: 包含报告名称和描述的报告对象 JSON 数组。
    """
    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:
    """加载报告的列定义,以便构建查询。

    参数:
        report_name: 报告名称(例如 'TradeActivity')。
        query_intent: 要提取数据的描述。

    返回:
        str: 包含参数和列定义的 JSON 对象。
    """
    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]:
    """运行预定义报告。该工具会验证每个过滤条件与报告模式的匹配情况,并构建参数化 SQL 查询。LLM 从不编写原始 SQL,因此过滤值无法注入查询。

    参数:
        report_name: 来自 `get_report_list` 的报告(例如 'TradeActivity')。
        filters: 以列名为键的等值过滤器,例如
            {"symbol": "AAPL", "date": "2024-03-15"}。
        limit: 可选的行数上限(1..10000)。

    返回:
        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"]}

    # 拒绝任何不在报告允许列列表中的过滤字段。
    unknown = set(filters) - allowed_columns
    if unknown:
        raise ValueError(
            f"未知的过滤字段 {sorted(unknown)} for {report_name}。 "
            f"允许的字段: {sorted(allowed_columns)}"
        )

    # 使用命名绑定参数构建 SQL。
    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 必须是 [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:宏观工作流编排

LangGraph 通过三大核心能力为多智能体系统提供生产级编排,这使其非常适合处理复杂的 AI 工作流。

基于图的状态机:LangGraph 将代理工作流建模为有向图,其中节点表示包含代理逻辑的函数,边决定执行流程。这种声明式方法将复杂的多步骤推理转化为可读且易于维护的代码。图结构支持条件分支、并行执行和动态路由。这些能力对于需要根据中间结果动态调整工作流的真实场景至关重要。

持久化状态管理:框架的检查点系统会在每个节点执行后自动对整个工作流状态进行快照记录。通过这些检查点,您可以平滑地从故障中恢复,并支持人机协作交互。当分析师需要查看中间结果或发生错误时,系统会从精确的检查点恢复,然后继续执行而不会丢失之前的工作。这种有状态的架构支持多轮对话、迭代优化和跨数小时或数天的长期调查等常见模式。

生产可靠性:LangGraph 内置了具有指数退避机制的重试策略,用于处理系统中可能出现的限流或其他故障。它还通过 OpenTelemetry 提供全面的可观测性,便于与大多数可观测性应用进行集成。

让我们构建一个编排层,用于在我们的专业 Strands 代理之间路由查询。以下代码定义了我们的工作流图:共享状态、选择调用哪些专业代理的调度器、它们之间的条件路由,以及基于检查点的持久化以实现恢复和人机协作审查。

code
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]
    # 每个专业代理在此处写入其见解
    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)

code

市场分析代理的 LangGraph 工作流编排

## 为什么选择 LangGraph 与 Strands 结合

许多企业用例必须基于严格预定义的工作流。单纯依赖大语言模型的非确定性来执行正确步骤存在风险。然而,这些工作流中的特定步骤确实需要大语言模型的敏捷推理能力。

LangGraph 与 Strands 的结合弥合了这一差距。您可以使用它构建系统,其中确定性编排包含本地化的动态智能。

这种架构组合如何解决复杂的工作流挑战:

**节点级智能**:LangGraph 定义工作流的高层编排。Strands 代理被放置在特定节点中,复杂工作流可能需要大语言模型分析或处理模糊性。它们仅在严格需要灵活性的地方应用自主推理和工具使用。

**节点代理的上下文隔离**:单体代理很容易丢失指令。通过将独立的 Strands 代理放置在 LangGraph 的不同节点中,您可以隔离内存。每个代理管理自己的高度专注的上下文和工具历史,而 LangGraph 保持结构化的总体会话状态,这些不同代理可以独立更新和引用。

**增强编排能力**:本质上,LangGraph 是一个底层路由和编排工具。通过嵌入 Strands,您可以立即将全面的企业代理框架插入到您的图中。您将获得 LangGraph 的强大路由能力,以及 Strands 本地的模型上下文协议(MCP)集成、控制转向、安全防护和评估功能。

具体来说,每个专家都是一个 LangGraph 节点。该节点启动一个带有独立系统提示、工具和隔离上下文的新 Strands 代理。它执行协调器分配的任务并返回状态更新。LangGraph 将这些部分更新合并到其他节点读取的共享状态中。以下是安全监控节点的示例:

async def security_monitor_node(state: AgentState) -> AgentState: """ 单日活动分析师代理,用于评估价格、成交量和逐笔交易 """ 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, )

从协调器填充的共享状态中获取该节点的任务

task = state.get("agent_task_map", {}).get("security_monitor", state["query_text"])

code

Strands 运行自己的推理 + 工具循环;收集代理的最终文本

chunks = [] async for event in agent.stream_async(task): if "data" in event: chunks.append(event["data"]) result = "".join(chunks)

返回共享状态更新

return { "security_monitor_insights": {"task": task, "business_insights": result}, "current_agent_index": state.get("current_agent_index", 0) + 1, }

code

## 使用 AgentCore 的 AWS 基础设施和部署

Amazon Bedrock AgentCore 提供了全托管服务,用于大规模部署和运行代理,既能减轻基础设施管理负担,又能提供生产级功能。

运行时部署:作为 Amazon Bedrock AgentCore 的一项功能,AgentCore 运行时通过最小的配置即可将本地代理代码转换为云原生部署。该服务与框架无关,开箱即用支持 LangGraph 和 Strands。它为动态代理工作负载提供专用基础设施,包括用于长期调查的扩展运行时、用于交互式工作流的低延迟执行,以及基于需求的自动扩展。

使用 AgentCore Python SDK(Amazon Bedrock AgentCore 的一项功能)和入门工具包部署您的 LangGraph 编排器和 Strands 代理。运行时会自动处理容器化、网络和计算资源分配。AgentCore 负责容器编排、扩展和会话管理等基础性繁重工作。

api.py --- AgentCore 运行时入口点

from bedrock_agentcore.runtime import BedrockAgentCoreApp from src.agents import Workflow

app = BedrockAgentCoreApp() workflow = Workflow()

@app.entrypoint async def market_surveillance_workflow(payload): """由 AgentCore 在每次请求时调用。生成流式数据块。""" 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()

使用 AgentCore 入门工具包进行部署

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"代理 ARN: {result.agent_arn}")

调用已部署的代理

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": "2024 年 3 月 15 日上午 11:00 AAPL 股价飙升的原因是什么?", "session_id": "session-001", "actor_id": "analyst-jane", }), )

code

# AgentCore 返回一个服务器发送事件流。解析方式如下:
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  # 数据块格式错误 --- 跳过,不中断
    if isinstance(chunk, dict) and chunk.get("type") == "text":
        print(chunk["content"], end="")

内存集成:LangGraph 通过 langgraph-checkpoint-aws 包与 AgentCore 内存集成(这是 Amazon Bedrock AgentCore 的一项功能),仅需几行代码即可实现短期检查点持久化和智能长期记忆检索。

本示例中使用的 AgentCoreMemorySaver 类负责处理包含用户消息、AI 响应、图执行状态和元数据的检查点对象。每个节点执行完成后,LangGraph 会自动将检查点保存至 AgentCore 内存。无需管理 Amazon DynamoDB 表或实现自定义序列化逻辑,即可获得有状态的对话和工作流恢复能力。

AgentCoreMemoryStore 类提供智能记忆功能,AgentCore 会自动从对话中提取洞察、摘要和用户偏好。代理可以在未来的交互中搜索这些记忆,从而提供随时间改进的个性化体验。这解决了代理无状态的根本难题:每次交互都基于先前知识进行,而非从零开始。

code
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="用于市场监控多代理工作流的记忆。",
    eventExpiryDuration=90,  # 天
)
MEMORY_ID = response["memory"]["id"]
print(f"记忆 ID: {MEMORY_ID}")

# 等待状态变为 ACTIVE,设定 10 分钟截止时间。创建通常需要 1-3 分钟。
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_ID} 状态为 {status!r}(预期为 ACTIVE)")
    time.sleep(10)

构建图时:

code
checkpointer = AgentCoreMemorySaver(MEMORY_ID, region_name=REGION)
graph = workflow.compile(checkpointer=checkpointer)

调用时,向代理传递 thread_id 和 actor_id。这些是用户和会话的唯一标识符:

code
config = {
    "configurable": {
        "thread_id": "surveillance-session-001",
        "actor_id": "analyst-jane",
    }
}

response = await graph.ainvoke(
    {"query_text": "哪些经纪人最活跃?"},
    config=config,
)

可观测性与运维:AgentCore 通过与 Amazon CloudWatch 和 AWS X-Ray 的集成提供内置可观测性,可捕获代理执行轨迹、工具调用和性能指标。该服务提供仪表板用于监控代理行为、识别瓶颈和优化成本。结合 LangGraph 的 OpenTelemetry 事件,您可以从高层工作流编排到单个 LLM 调用和推理步骤获得全面的可视性。

结论

在本文中,我们展示了如何通过将 LangGraph 强大的工作流编排与 Strands 的智能代理推理能力相结合,并部署在 AWS 基础设施上,构建生产就绪的多智能体 AI 系统。

我们探讨的混合架构证明了 LangGraph 在宏观层面编排(管理代理协调、状态持久化和工作流恢复)方面的优势,而 Strands 则在单个节点内提供详细的推理引擎。通过这种职责分离,您可以构建能够处理复杂业务流程的复杂系统,同时通过基于检查点的恢复和全面的可观测性保持生产可靠性。

可以将这种架构扩展到其他复杂的编排场景,例如文档处理流水线、客户服务自动化或合规性监控系统。Strands 的模型无关性结合 LangGraph 的状态管理,使这种模式对于需要灵活性和可靠性的企业应用尤其有价值。

若想深入了解技术实现细节并亲自构建市场分析代理,请查看 GitHub 仓库 。

关于作者

'"`