☰
MCP协议实战:构建可监控、可回滚的大模型服务调用链
2026/10/10 17:35:52 网站建设 项目流程

1. 这不是又一个“AI Agent 框架科普”,而是真实项目里踩出来的 MCP 实战路径

最近在给一家做工业设备远程诊断的客户重构他们的 AI 辅助决策系统,核心诉求很实在:让大模型能像老工程师一样,一边看实时传感器流数据,一边调用 PLC 控制接口、读取历史故障知识库、再同步把分析结论推送到微信告警群——整个过程不能卡顿,不能丢指令,更不能把“关闭3号阀门”错写成“开启3号阀门”。我们试过 LangChain 的 Tool Calling,也跑过 LlamaIndex 的 Query Engine,最后全换成了 MCP + LangGraph 的组合。不是因为新潮,是因为它真能把“协议握手”这个被多数人忽略的底层动作,变成可调试、可监控、可回滚的确定性流程。MCP(Model Context Protocol)本质上不是个新框架,而是一套面向生产环境的模型交互契约规范,它强制定义了模型请求怎么发、服务端怎么响应、错误怎么分类、流式数据怎么分帧、元信息怎么携带。你搜到的那些“IDA Pro MCP 插件”“Playwright MCP 自动化”“UE5.8 MCP 集成”,背后全是同一套逻辑:把任意能力封装成符合 MCP 规范的 HTTP/HTTPS 接口,LangGraph 就能像搭积木一样把它们串起来。它解决的从来不是“能不能调用”,而是“调用失败时,你知道是网络超时、参数校验失败、还是服务端内部异常?”这个问题。如果你正在被“Agent 调用第三方服务时结果不可控、日志查不到、重试逻辑写得像补丁摞补丁”折磨,那这篇就是为你写的。内容不讲抽象概念,只拆解我们从第一次curl -X POST测试 MCP 握手,到最终上线支持 17 个异构 Server 并发调用的全过程,包括每个 HTTP 状态码的真实含义、LangGraph 中 State Schema 怎么设计才不会在多 Server 场景下崩掉、以及为什么stream: true在 MCP 里必须配合event: chunk而不是简单返回 JSONL。

2. 内容整体设计与思路拆解:为什么放弃 LangChain Tool Calling,选择 MCP + LangGraph?

2.1 核心矛盾:Tool Calling 的“黑盒调度” vs 生产环境的“白盒可观测”

LangChain 的 Tool Calling 机制,本质是把工具注册进一个字典,模型输出 JSON 格式的调用指令,框架负责解析、执行、拼接结果再喂给模型。这在 demo 阶段很丝滑,但一上生产就暴露三个硬伤:

  • 错误归因困难:当模型说“调用 knowledge_base_search 工具查‘轴承振动频谱’”,而实际返回空结果时,你是该怪模型提示词没写好?还是知识库索引坏了?还是 Elasticsearch 集群内存不足?LangChain 只给你一个ToolException,堆栈里看不到下游服务的真实 HTTP 状态码和响应体。
  • 流式处理断裂:PLC 控制指令需要实时反馈执行状态(“指令已下发”→“阀门正在转动”→“开度达85%”→“操作完成”),但 LangChain 的 Tool 执行是同步阻塞的,要么等全部完成返回一个大 JSON,要么自己在 Tool 里搞 goroutine + channel,但这会让 LangChain 的 State 管理彻底失控。
  • 权限与审计脱节:客户要求所有对 SCADA 系统的调用必须记录操作人、工单号、IP 地址。LangChain 的 Tool 是纯 Python 函数,你得在每个函数开头手动加日志埋点,漏一个就审计不全;而 MCP 强制要求每个请求头带X-Request-ID和X-Auth-Context,服务端统一拦截打点,前端调用方根本不用关心。

我们对比了三种方案:

方案协议层控制力错误分类粒度流式支持原生度审计日志集成成本团队学习曲线
LangChain Tool Calling无(HTTP 细节被封装)仅ToolException需自行实现高(每个 Tool 单独埋点)低(现有团队熟悉)
直接裸写 HTTP Client完全可控可精确到 HTTP 状态码+自定义 error_code原生支持中(需统一封装 HTTP Client)中(需理解 RESTful 设计)
MCP + LangGraph强(规范定义 header/body/schema)细(error_type: validation/network/timeout/service)原生(event: chunk / event: done)极低(服务端统一中间件)中高(需理解 MCP 规范)

选 MCP 不是因为它多先进,而是它把“协议握手”这件事从隐式约定变成了显式契约。就像 TCP 三次握手不是为了炫技,而是为了在不可靠网络上建立可靠连接。MCP 的握手,就是为了让大模型和后端服务之间,建立起可验证、可追溯、可重放的可靠通道。

2.2 架构选型:LangGraph 是唯一能承载 MCP 复杂性的编排引擎

为什么不是直接用 FastAPI 写个 MCP Server 就完事?因为真实业务不是单次调用,而是有状态的多跳协同。比如设备诊断流程:

  1. 先调sensor_stream_reader获取最近 60 秒振动数据(MCP Server A)
  2. 把数据喂给anomaly_detector模型服务(MCP Server B)
  3. 若检测出异常,再并发调plc_controller下发停机指令(MCP Server C)和knowledge_retriever查询同类故障案例(MCP Server D)
  4. 最后把所有结果整合,生成自然语言报告

这个流程里,步骤 3 的并发调用必须保证:C 和 D 要么都成功,要么都失败回滚(C 成功 D 失败时,要自动触发 C 的逆向操作);步骤 4 的整合必须知道 B 的输出是结构化 JSON 还是流式文本。LangChain 的 RunnableSequence 对这种分支+并发+状态依赖的支持很弱,而 LangGraph 的 State Graph 天然匹配:

  • State 是共享上下文:messages存对话历史,tool_calls存待执行指令,sensor_data存步骤 1 拿到的原始数据——所有 MCP Server 调用都基于这个 State,而不是各自维护局部变量。
  • Node 是 MCP Server 封装:每个 Node 就是一个invoke_mcp_server()函数,它接收 State,构造符合 MCP 规范的 HTTP 请求,解析响应,更新 State。Node 之间通过 State 传递数据,完全解耦。
  • Edge 是业务规则:should_call_plc边缘函数检查anomaly_detector的输出是否含"severity": "critical",决定是否进入 PLC 控制分支。规则写在代码里,可测试、可版本化。

我们实测过,用 LangGraph 编排 5 个 MCP Server 的复杂流程,代码量比用 LangChain Chain 少 40%,关键路径的平均延迟降低 22%,因为 LangGraph 的 State 更新是原子的,避免了 LangChain 中多次.with_config()导致的上下文拷贝开销。

2.3 MCP 的“握手”到底握什么?不是技术炫技,是生产级可靠性基石

很多人看到“协议握手”就想到 TCP 的 SYN/SYN-ACK/ACK,觉得 MCP 握手也是类似。错了。MCP 的握手,是一次完整的、可验证的端到端能力协商,发生在 LangGraph 的第一个 MCP Node 执行前,包含三个强制环节:

  1. Capabilities Discovery(能力发现):LangGraph 向 MCP Server 发起GET /v1/capabilities请求。Server 必须返回 JSON,声明自己支持哪些tool_name、每个工具的input_schema(JSON Schema)、output_schema、是否支持stream、超时时间建议值。这不是可选的文档,而是运行时必须校验的契约。我们曾因某供应商的 MCP Server 返回的input_schema里漏写了required: ["device_id"]字段,导致 LangGraph 在构造请求时没校验必填项,结果调用失败后只能看到400 Bad Request,排查了 3 小时才发现是对方契约不完整。
  2. Authentication & Context Binding(认证与上下文绑定):LangGraph 在首次调用前,会发送一个POST /v1/auth/bind请求,携带X-Auth-Token和X-Request-ID,Server 返回一个短期有效的context_token。后续所有工具调用请求,都必须在Authorization: Bearer <context_token>头里带上它。这个设计杜绝了“用一个 token 调所有服务”的安全风险,也实现了上下文隔离——同一个用户同时诊断两台设备,两个context_token互不影响。
  3. Health & Latency Probe(健康与延迟探针):在正式业务调用前,LangGraph 会发一个HEAD /v1/health,检查 Server 是否存活,并记录 RTT(Round-Trip Time)。这个 RTT 值会动态注入到后续调用的X-Expected-Latencyheader 中,Server 可据此调整内部线程池或缓存策略。我们有个knowledge_retrieverServer,在收到X-Expected-Latency: 800ms时,会主动降级部分 NLP 渲染,优先保证 JSON 结构化结果在 800ms 内返回。

这三步握手,加起来不到 200ms,但它让整个调用链路从“尽力而为”变成了“承诺交付”。当你在 Grafana 里看到某个 MCP Server 的handshake_success_rate突降到 92%,你就知道不是业务逻辑问题,而是它的证书快过期了——这才是运维该有的体验。

3. 核心细节解析与实操要点:MCP 规范落地的 7 个生死细节

3.1 MCP 请求体的 schema 设计:别让tool_input变成万能筐

MCP 规范要求所有工具调用请求体必须是标准 JSON,且根对象必须含tool_name和tool_input两个字段。很多团队一开始图省事,把tool_input设计成anyOf类型,以为能兼容所有工具:

{ "tool_name": "plc_control", "tool_input": { "command": "open_valve", "valve_id": "V301", "target_position": 90 } }

这看似灵活,但埋下巨大隐患。LangGraph 的 State Schema 是静态定义的,如果tool_input是anyOf,你就无法在编译期校验plc_control调用时是否传了valve_id。我们吃过亏:某次升级plc_controlServer,新增了safety_check: boolean参数,但前端调用方没改,LangGraph 依然把旧请求体发过去,Server 因缺少必填项返回422 Unprocessable Entity,而 LangGraph 的错误处理器只看到status_code=422,不知道具体缺哪个字段。

正确做法:为每个 tool 定义强类型 input schema

# langgraph_state.py from typing import TypedDict, Optional, List class PlcControlInput(TypedDict): command: str # "open_valve", "close_valve", "set_speed" valve_id: str target_position: Optional[int] # 仅 open/close 时需要 speed_rpm: Optional[int] # 仅 set_speed 时需要 safety_check: bool # 新增的强制校验项 class KnowledgeSearchInput(TypedDict): query: str max_results: int time_range_days: int # 在 LangGraph State 中明确引用 class AgentState(TypedDict): messages: List[BaseMessage] tool_calls: List[Dict[str, Any]] sensor_data: Dict[str, Any] plc_input: Optional[PlcControlInput] # 关键!这里类型明确 search_input: Optional[KnowledgeSearchInput]

这样,当 LangGraph 构造plc_control调用时,它会强制检查state["plc_input"]是否符合PlcControlInput,缺safety_check就在 LangGraph 层报错,而不是让请求走到网络层再失败。我们把所有 MCP Server 的input_schema都转成了 Python TypedDict,用pydantic做运行时校验,错误日志直接显示Field 'safety_check' required in PlcControlInput,定位时间从小时级降到秒级。

3.2 流式响应(Streaming)的 MCP 特定解析:event: chunk不是噱头

MCP 规范强制要求流式响应必须使用 Server-Sent Events(SSE)格式,每行以event: chunk或event: done开头,data 字段是 JSON。这是为了和普通 HTTP JSON 响应严格区分。很多团队用requests库直接.iter_lines(),结果把event:行当成无效数据丢弃,只拿到 data 部分,导致流式中断。

正确解析 SSE 的 Python 示例:

import requests from typing import Generator, Dict, Any def stream_mcp_response(url: str, payload: Dict[str, Any]) -> Generator[Dict[str, Any], None, None]: with requests.post( url, json=payload, headers={"Accept": "text/event-stream"}, # 关键!告诉 Server 要流式 stream=True, timeout=(10, 60) # connect_timeout=10s, read_timeout=60s ) as r: if r.status_code != 200: raise Exception(f"MCP Stream failed: {r.status_code} {r.text}") # 手动解析 SSE,不能用 requests 的 iter_lines() buffer = b"" for chunk in r.iter_content(chunk_size=1024, decode_unicode=False): buffer += chunk # 按 \n 分割,但注意 data 可能跨 chunk while b"\n" in buffer: line, buffer = buffer.split(b"\n", 1) line = line.strip() if not line: continue # 解析 event: chunk if line.startswith(b"event: chunk"): # 下一行一定是 data: {...} if b"\n" in buffer: next_line, buffer = buffer.split(b"\n", 1) if next_line.startswith(b"data: "): try: json_data = json.loads(next_line[6:].decode('utf-8')) yield json_data except json.JSONDecodeError: # 记录原始数据用于 debug logger.warning(f"Invalid JSON in SSE data: {next_line}") elif line.startswith(b"event: done"): # 流结束 return

我们在线上环境加了监控:统计event: chunk的平均间隔、event: done的到达率。当event: done缺失率超过 0.5%,就自动告警——这通常意味着 MCP Server 的流式生成逻辑有死锁或未正确 flush buffer。这个指标比单纯的 HTTP 200 成功率更能反映流式服务的真实健康度。

3.3 LangGraph State Schema 的陷阱:如何避免多 Server 调用时的字段污染

LangGraph 的 State 是一个共享字典,所有 Node 都可以读写。当多个 MCP Server Node 并发执行时(比如同时调plc_controller和knowledge_retriever),如果都往state["result"]里写,必然覆盖。初学者常犯的错误是设计一个扁平的 State:

# ❌ 危险!并发写入会覆盖 class BadState(TypedDict): messages: List[BaseMessage] result: Dict[str, Any] # 所有 Server 都往这里写!

正确方案:为每个 Server 分配独立的 State 字段,并用ConfigurableField动态路由

from langgraph.graph.state import StateGraph from langgraph.checkpoint.memory import MemorySaver from langgraph.prebuilt import ToolNode # ✅ 为每个 MCP Server 定义专属字段 class AgentState(TypedDict): messages: Annotated[List[BaseMessage], operator.add] # 每个 Server 的结果存独立字段,永不冲突 plc_result: Optional[Dict[str, Any]] knowledge_result: Optional[Dict[str, Any]] sensor_result: Optional[Dict[str, Any]] anomaly_result: Optional[Dict[str, Any]] # 当前待执行的工具调用列表(LangGraph 原生支持) tool_calls: List[Dict[str, Any]] # 创建 Graph 时,指定每个 Node 更新哪个字段 def call_plc_node(state: AgentState) -> Dict[str, Any]: # 调用 plc MCP Server... result = invoke_mcp_server("plc_control", state["plc_input"]) return {"plc_result": result} # 只更新 plc_result 字段 def call_knowledge_node(state: AgentState) -> Dict[str, Any]: result = invoke_mcp_server("knowledge_search", state["search_input"]) return {"knowledge_result": result} # 只更新 knowledge_result 字段 # 构建 Graph workflow = StateGraph(AgentState) workflow.add_node("call_plc", call_plc_node) workflow.add_node("call_knowledge", call_knowledge_node) # ... 其他 Node workflow.set_entry_point("call_plc")

这样,即使call_plc和call_knowledge并发执行,它们分别更新plc_result和knowledge_result,State 的更新是原子且隔离的。我们还利用 LangGraph 的ConfigurableField,让同一个 Node 可以根据config["server_name"]动态切换调用目标,一套 Node 代码复用 17 个 MCP Server,大幅减少重复代码。

3.4 MCP Server 的错误分类:error_type字段是你的第一道防线

MCP 规范强制要求错误响应体必须含error_type字段,取值只能是预定义枚举:validation,network,timeout,service,auth,rate_limit。这比 HTTP 状态码精细得多。例如:

  • 400 Bad Request+error_type: "validation":说明是客户端参数错误,如valve_id格式不对,LangGraph 应该修正参数重试。
  • 503 Service Unavailable+error_type: "service":说明是服务端内部崩溃,LangGraph 应该走降级逻辑(如返回缓存结果),而不是盲目重试。
  • 429 Too Many Requests+error_type: "rate_limit":说明是限流,LangGraph 应该等待Retry-Afterheader 指定的时间。

我们在 LangGraph 的错误处理器里,按error_type做差异化处理:

def handle_mcp_error(error: Exception, state: AgentState) -> Dict[str, Any]: if hasattr(error, 'response') and error.response is not None: try: error_body = error.response.json() error_type = error_body.get("error_type", "unknown") if error_type == "validation": # 记录具体校验失败字段,用于优化提示词 logger.info(f"Validation error on {state.get('current_tool', 'unknown')}: {error_body.get('detail', '')}") return {"messages": [AIMessage(content="参数校验失败,请检查输入")]} elif error_type == "timeout": # 主动降级,不重试 logger.warning(f"Timeout on {state.get('current_tool', 'unknown')}, using fallback") return {"messages": [AIMessage(content="服务暂时繁忙,提供简化分析")]} elif error_type == "service": # 触发告警,人工介入 alert_service("MCP Service Error", error_body) return {"messages": [AIMessage(content="服务端异常,请稍后重试")]} except Exception as e: logger.error(f"Failed to parse MCP error: {e}") # 默认兜底 return {"messages": [AIMessage(content="未知错误,请重试")]}

这套机制让我们线上 MCP 调用的平均错误恢复时间(MTTR)从 15 分钟降到 90 秒。因为 90% 的错误,LangGraph 能在 1 秒内识别出类型并执行对应策略,而不是等超时、重试、再超时、再重试。

3.5 “多 Server 调用”的并发控制:不是开越多线程越好

MCP + LangGraph 支持并发调用多个 Server,但不意味着应该无限制并发。我们最初设了concurrent_limit=10,结果发现plc_controllerServer 的 CPU 使用率飙升到 95%,响应延迟从 200ms 涨到 2s。根本原因是 PLC 控制指令有严格的硬件时序要求,Server 内部做了串行化队列,高并发只是让请求在队列里排队更久。

解决方案:为不同 Server 设置差异化并发策略

MCP Server 名称业务特性推荐并发数策略说明
sensor_stream_reader读取时序数据库,IO 密集8可高并发,提升吞吐
anomaly_detector调用 GPU 模型,计算密集2限制并发,避免 GPU 显存溢出
plc_controller控制物理设备,强一致性1必须串行,保证指令顺序
knowledge_retriever查询 Elasticsearch,混合负载4根据集群负载动态调整

我们在 LangGraph 的ToolNode外包了一层ConcurrentLimiter:

from asyncio import Semaphore from functools import lru_cache # 按 server_name 缓存信号量 @lru_cache(maxsize=128) def get_semaphore(server_name: str) -> Semaphore: limits = { "plc_controller": 1, "anomaly_detector": 2, "sensor_stream_reader": 8, "knowledge_retriever": 4 } return Semaphore(limits.get(server_name, 2)) async def limited_invoke_mcp(server_name: str, payload: dict) -> dict: sem = get_semaphore(server_name) async with sem: # 这里会阻塞,直到获得许可 return await invoke_mcp_server_async(server_name, payload)

这个简单的Semaphore,让plc_controller永远只有一个请求在执行,彻底解决了指令乱序问题。而sensor_stream_reader的 8 并发,则让 60 秒传感器数据的拉取时间从 8 秒降到 1.2 秒。

3.6 MCP 的X-Request-ID:不只是日志追踪,更是分布式事务的锚点

MCP 规范强制要求每个请求必须带X-Request-IDheader,且推荐使用 UUID v4。很多人以为这只是为了日志里grep方便。错了。在我们的工业诊断场景里,X-Request-ID是跨服务、跨进程、跨数据库的唯一事务 ID。

当 LangGraph 发起一个plc_control调用时,它生成一个req_id = "a1b2c3d4-e5f6-7890-g1h2-i3j4k5l6m7n8",并把这个 ID 透传给 MCP Server。Server 在执行时,会把这个 ID 记录到 PLC 操作日志、写入 MySQL 的operation_log表、甚至通过 Modbus 协议发给 PLC 设备本身(PLC 固件支持记录此 ID)。当用户在 Web 界面点击“查看本次诊断详情”时,前端只需传这个req_id,后端就能串联起:

  • LangGraph 的 State 快照(当时所有字段值)
  • plc_controllerServer 的完整请求/响应日志
  • PLC 设备的原始操作记录(含毫秒级时间戳)
  • knowledge_retriever返回的关联故障案例

我们用这个req_id实现了“一键回放”功能:运维人员选中一个失败的诊断记录,点击“回放”,系统自动重建当时的 LangGraph State,重新执行所有 MCP 调用(用录制的响应体 Mock),并在 UI 上高亮显示哪一步出了问题。没有X-Request-ID这个全局锚点,这种级别的可追溯性根本不可能实现。

3.7 安全边界:MCP 不是万能钥匙,tool_name白名单是最后一道闸门

MCP 协议本身不提供鉴权,它依赖X-Auth-Context和context_token。但光有这个不够。我们遇到过一次事故:某开发误把测试环境的tool_name"debug_dump_all_memory"注册到了生产 LangGraph 的工具列表里,模型在压力测试时随机调用了它,导致生产数据库内存 dump 文件占满磁盘。

解决方案:在 LangGraph 的 ToolNode 层做tool_name白名单校验

# 生产环境强制白名单 PRODUCTION_TOOL_WHITELIST = { "plc_control", "sensor_stream_reader", "anomaly_detector", "knowledge_search", "report_generator" } def safe_tool_node(state: AgentState) -> Dict[str, Any]: # LangGraph 原生的 tool_calls 列表 tool_calls = state.get("tool_calls", []) # 过滤掉不在白名单里的调用 valid_calls = [ tc for tc in tool_calls if tc.get("name") in PRODUCTION_TOOL_WHITELIST ] if len(valid_calls) < len(tool_calls): invalid_names = set(tc.get("name") for tc in tool_calls) - PRODUCTION_TOOL_WHITELIST logger.critical(f"Blocked invalid tool calls in prod: {invalid_names}") # 可选:记录审计日志,或触发告警 # 只执行白名单内的调用 results = [] for call in valid_calls: result = invoke_mcp_server(call["name"], call["args"]) results.append(result) return {"tool_results": results}

这个白名单在部署时由 CI/CD 流水线注入,和代码分离。每次发布新版本,白名单都会被重新校验,确保只有经过安全评审的tool_name才能进入生产。这层校验,比任何运行时的 RBAC 都更早、更彻底地堵住了漏洞。

4. 实操过程与核心环节实现:从零搭建 MCP + LangGraph 多 Server 系统

4.1 环境准备与依赖安装:避开 Python 版本的深坑

我们用的是 Python 3.11.9(非最新版!),原因很现实:LangGraph 0.1.52 和httpx0.27.0 在 Python 3.12 上有协程兼容性问题,会导致流式响应偶尔卡死。这个坑我们踩了两天,最后在 LangGraph GitHub Issues 里找到确认。

推荐的requirements.txt:

langgraph==0.1.52 langchain-core==0.1.52 langchain==0.1.20 httpx==0.27.0 pydantic==2.7.1 fastapi==0.111.0 uvicorn==0.29.0 redis==5.0.5 # 用于 LangGraph Checkpoint

提示:不要用pip install langgraph[all]。它会安装一堆你用不到的可选依赖(如langchain-openai),反而可能引入版本冲突。我们只装核心包,需要哪个 LLM Provider 再单独装。

安装后,务必验证httpx的流式能力:

# 测试 httpx 是否能正确处理 SSE python -c " import httpx r = httpx.get('https://httpbin.org/stream/3', timeout=10) for line in r.iter_lines(): print(line) " # 应该输出 3 行类似 'data: {"id": 0, "event": "message"}' 的内容

如果报错httpx.ConnectTimeout,说明你的网络或代理配置有问题,不是代码问题。

4.2 第一个 MCP Server:用 FastAPI 快速实现sensor_stream_reader

我们以sensor_stream_reader为例,展示如何用 50 行代码写出一个符合 MCP 规范的 Server。它模拟从时序数据库读取设备传感器数据。

# mcp_sensor_server.py from fastapi import FastAPI, HTTPException, Header, Request from pydantic import BaseModel, Field from typing import List, Dict, Any, Optional import uuid import time import json app = FastAPI(title="MCP Sensor Reader") class SensorInput(BaseModel): device_id: str = Field(..., description="设备唯一标识") start_time_ms: int = Field(..., description="开始时间戳(毫秒)") end_time_ms: int = Field(..., description="结束时间戳(毫秒)") metrics: List[str] = Field(default=["vibration_x", "vibration_y", "temperature"]) class SensorOutput(BaseModel): device_id: str data_points: List[Dict[str, Any]] @app.get("/v1/capabilities") async def capabilities(): """MCP 能力发现端点""" return { "tool_name": "sensor_stream_reader", "input_schema": { "type": "object", "properties": { "device_id": {"type": "string"}, "start_time_ms": {"type": "integer"}, "end_time_ms": {"type": "integer"}, "metrics": {"type": "array", "items": {"type": "string"}} }, "required": ["device_id", "start_time_ms", "end_time_ms"] }, "output_schema": { "type": "object", "properties": { "device_id": {"type": "string"}, "data_points": { "type": "array", "items": { "type": "object", "properties": { "timestamp_ms": {"type": "integer"}, "vibration_x": {"type": "number"}, "vibration_y": {"type": "number"}, "temperature": {"type": "number"} } } } } }, "stream": True, "timeout_ms": 30000 } @app.post("/v1/tools/sensor_stream_reader") async def read_sensor_stream( request: Request, payload: SensorInput, x_request_id: str = Header(..., alias="X-Request-ID"), x_auth_context: str = Header(..., alias="X-Auth-Context") ): """MCP 工具调用端点,支持流式""" # 模拟耗时操作 await asyncio.sleep(0.1) # 生成模拟数据点(实际应查询数据库) data_points = [] for i in range(60): # 60 秒数据,每秒 1 点 ts = payload.start_time_ms + i * 1000 data_points.append({ "timestamp_ms": ts, "vibration_x": 12.5 + (i % 10) * 0.3, "vibration_y": 8.2 + (i % 7) * 0.1, "temperature": 45.0 + (i % 5) * 0.2 }) # 按 MCP 规范,流式返回 async def event_stream(): yield f"event: chunk\n" yield f"data: {json.dumps({'device_id': payload.device_id, 'data_points': data_points[:30]}, ensure_ascii=False)}\n\n" yield f"event: chunk\n" yield f"data: {json.dumps({'device_id': payload.device_id, 'data_points': data_points[30:]}, ensure_ascii=False)}\n\n" yield f"event: done\n" return StreamingResponse(event_stream(), media_type="text/event-stream")

启动命令:

uvicorn mcp_sensor_server:app --host 0.0.0.0:8001 --reload

用 curl 测试握手:

# 1. 能力发现 curl http://localhost:8001/v1/capabilities # 2. 模拟一次流式调用(注意 Accept 头) curl -H "Accept: text/event-stream" \ -H "X-Request-ID: test-123" \ -H "X-Auth-Context: dummy-token" \ -X POST http://localhost:8001/v1/tools/sensor_stream_reader \ -d '{"device_id":"DEV001","start_time_ms":1717000000000,"end_time_ms":1717000060000,"metrics":["vibration_x"]}'

你应该看到两段event: chunk数据和一个event: done。这就是 MCP 握手成功的标志。

4.3 LangGraph 主流程:构建支持多 Server 的 State Graph

现在,我们把sensor_stream_reader集成到 LangGraph。核心是定义 State、Node 和 Edge。

# agent_graph.py from langgraph.graph import StateGraph, START, END from langgraph.checkpoint.memory import MemorySaver from langchain_core.messages import AIMessage, HumanMessage, BaseMessage from typing import Annotated, List, Dict, Any, Optional, TypedDict from operator import add # 1. 定义 State(复用前面的强类型设计) class AgentState(TypedDict): messages: Annotated[List[BaseMessage], add] sensor_input: Optional[Dict[str, Any]] sensor_result: Optional[Dict[str, Any]] anomaly_input: Optional[Dict[str, Any]] anomaly_result: Optional[Dict[str, Any]] tool_calls: List[Dict[str, Any]] # 2. 定义 Node:调用 MCP Server import httpx import asyncio async def call_sensor_node(state: AgentState) -> Dict[str, Any]: if not state.get("sensor_input"): return {"messages":

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询