用 LangGraph、Strands 与 AgentCore 构建可恢复的市场监控多智能体系统

2026-07-29 14 预计阅读时间: 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 分钟

市场监控不是一次简单的模型调用。一个生产系统需要持续接收行情与新闻,让不同智能体分别判断异常、核查证据、评估风险,并在任务中断后从最近状态继续执行。LangGraph、Strands 与 Amazon Bedrock AgentCore 的组合,正好对应这三个层次:LangGraph 管理确定性的工作流,Strands 承担智能体推理,AgentCore 提供运行时、记忆与可观测能力。

把编排与推理解耦

多智能体系统最容易失控的地方,是把业务流程全部塞进提示词。例如,让一个智能体自行决定是否调用其他智能体、失败后从哪里重试、结果应该保存到哪里。这样做虽然原型快,但执行路径难以测试,故障也很难定位。

更稳妥的架构是把职责分成三层:

  • LangGraph 工作流层:定义节点、状态、条件分支、重试边界与终止条件。
  • Strands 推理层:在单个节点内分析新闻、价格异动或风险信号,并调用被允许的工具。
  • AgentCore 平台层:承载智能体运行时,保存跨会话记忆,并收集调用链、延迟、错误和模型使用情况。

在市场监控场景中,一条典型路径可以表示为:

采集市场事件 -> 异常筛选 -> 新闻与公告核查 -> 风险评估 -> 告警或归档

关键点是让图决定“下一步做什么”,让智能体负责“这一步如何判断”。价格涨跌幅、数据是否齐全、风险分数是否超过阈值等规则,适合写成普通代码;语义归因、证据摘要和不确定性分析,则适合交给模型。

状态驱动让执行过程可检查

LangGraph 的状态不应只保存最终文本。生产环境中应保留能够解释决策、恢复任务和支持审计的结构化字段,例如:

  • event_id:事件的幂等键,避免同一行情事件产生重复告警。
  • symbolprice_change_pct:触发分析的原始市场数据。
  • evidence:新闻、公告或其他证据的摘要及引用标识。
  • risk_score:可用于条件路由的标准化分数。
  • decision:告警、观察或归档。
  • errors:节点错误及重试信息。

节点应尽量返回局部状态更新,而不是直接修改外部全局对象。这样既便于测试,也便于检查点系统在每个阶段保存一致快照。

下面是一个可以本地运行并继续改造的最小示例。它用 LangGraph 编排流程,用一个可替换的推理函数模拟 Strands 节点。安装依赖:

python -m venv .venv
source .venv/bin/activate
pip install langgraph

创建 market_graph.py

from typing import Literal, TypedDict

from langgraph.checkpoint.memory import MemorySaver
from langgraph.graph import END, START, StateGraph


class MarketState(TypedDict, total=False):
    event_id: str
    symbol: str
    price_change_pct: float
    evidence: list[str]
    risk_score: int
    decision: Literal["alert", "watch", "archive"]


def screen_event(state: MarketState) -> MarketState:
    change = abs(state["price_change_pct"])
    if change < 3.0:
        return {"risk_score": 10, "decision": "archive"}
    return {"risk_score": min(60, int(change * 8))}


def route_after_screen(state: MarketState) -> str:
    return "finish" if state.get("decision") == "archive" else "investigate"


def investigate(state: MarketState) -> MarketState:
    # 实际部署时,可将这里替换为 Strands Agent 调用,并为其配置
    # 新闻检索、公司公告查询等受控工具。
    evidence = [
        f"Observed an absolute move of {abs(state['price_change_pct']):.1f}%",
        "No verified corporate announcement was supplied to this demo",
    ]
    return {"evidence": evidence}


def assess_risk(state: MarketState) -> MarketState:
    score = state.get("risk_score", 0)
    if not state.get("evidence"):
        score += 20
    decision: Literal["alert", "watch", "archive"]
    decision = "alert" if score >= 70 else "watch"
    return {"risk_score": min(score, 100), "decision": decision}


builder = StateGraph(MarketState)
builder.add_node("screen", screen_event)
builder.add_node("investigate", investigate)
builder.add_node("assess", assess_risk)
builder.add_edge(START, "screen")
builder.add_conditional_edges(
    "screen",
    route_after_screen,
    {"finish": END, "investigate": "investigate"},
)
builder.add_edge("investigate", "assess")
builder.add_edge("assess", END)

graph = builder.compile(checkpointer=MemorySaver())

if __name__ == "__main__":
    config = {"configurable": {"thread_id": "event-NVDA-2025-001"}}
    result = graph.invoke(
        {
            "event_id": "NVDA-2025-001",
            "symbol": "NVDA",
            "price_change_pct": -8.4,
        },
        config=config,
    )
    print(result)

运行:

python market_graph.py

MemorySaver 适合本地演示,不适合多副本生产部署。接入 AgentCore 时,可以保留相同的状态模型和节点边界,再将临时检查点换成持久化后端,并把跨会话偏好、历史调查摘要等长期信息交给 AgentCore memory。具体 SDK 初始化参数应以当前 AgentCore 与 Strands 版本为准。

检查点恢复不等于重新执行全部任务

市场数据和外部工具调用通常具有时效性与成本。一个调查节点已经抓取过公告,后续风险评估失败时,不应默认重新抓取全部数据。检查点应在重要节点之后记录状态,使任务能够从最近成功阶段恢复。

恢复设计需要同时处理三个问题:

  1. 幂等性:以 event_id 加节点名构造写入键,重复执行不能生成多条告警。
  2. 数据新鲜度:恢复前检查行情和证据是否过期;过期时只回退到必要的采集节点。
  3. 副作用隔离:把发送告警、创建工单等动作放在独立节点,并记录外部系统返回的操作 ID。

长期记忆与检查点也需要区分。检查点回答“这个任务执行到了哪里”,长期记忆回答“过去对该公司、行业或用户形成了哪些可复用信息”。若把二者混在一起,任务状态会不断膨胀,过期结论也可能污染新的判断。

在 AgentCore 上关注可观测性边界

多智能体系统的日志不能只记录最终告警。一次错误结论可能来自数据源过期、工具超时、提示词变化、模型输出解析失败,或者路由条件本身写错。建议至少为每次执行记录:

  • 工作流运行 ID、事件 ID、线程 ID和节点名称。
  • 模型与提示词版本、工具调用名称和耗时。
  • 节点输入输出的结构摘要,而非无限制保存完整敏感内容。
  • 重试次数、异常类型和检查点恢复位置。
  • 最终风险分数、决策以及人工复核结果。

AgentCore 的 observability 能力可以作为统一追踪入口,但应用仍需主动传播运行 ID,并在 LangGraph 节点与 Strands 工具调用之间建立关联。对于行情、未公开公告或客户持仓等敏感数据,应配置脱敏、访问控制和保留周期,避免为了调试而永久存储完整提示词。

上线前的取舍清单

这套架构适合流程较长、需要恢复和审计的智能体任务,但也会引入更多状态管理与部署成本。若任务只是一次无副作用的摘要调用,完整多智能体图可能没有必要。

正式上线前,可以逐项确认:

  • 规则判断是否留在确定性代码中,而不是全部交给模型。
  • 每个节点是否具有明确输入、结构化输出和超时限制。
  • 检查点后端是否支持并发、持久化和预期的数据保留周期。
  • 告警节点是否幂等,恢复执行是否会重复触发外部动作。
  • 长期记忆是否设置来源、时间戳、有效期和删除机制。
  • 可观测数据是否足以重放决策,同时满足隐私与合规要求。
  • 高风险告警是否保留人工复核入口。

真正生产可用的多智能体系统,核心不在于部署了多少智能体,而在于每次判断能否被恢复、解释、追踪和约束。LangGraph 提供明确的执行骨架,Strands 封装节点内的推理能力,AgentCore 则补齐运行时、记忆和观测基础设施;三者的边界越清楚,系统越容易长期维护。


相关推荐