知识库查询 —— 网络搜索节点
本文档详细介绍知识库查询流程中的网络搜索节点(web_search_mcp)的设计与实现。该节点通过 MCP 协议调用外部网络搜索服务,获取实时的互联网搜索结果,作为本地知识库的补充信息源。
学习理念:网络搜索是多路检索的"兜底"通道——本地知识库查不到的、过时的、超出范围的内容,都靠它来补充。WebSearchMcpNode 通过 MCP 协议(Model Context Protocol)连接百炼等 AI 搜索服务,返回结构化搜索结果。核心价值:让 RAG 系统从"只查本地"升级为"本地 + 实时互联网"混合。
海外对标:MCP 协议由 Anthropic 于 2024 年底提出,对标 Google 的 Agent-to-Agent Protocol(A2A)和 OpenAI 的 Function Calling。
MCPServerSse+call_tool模式对标 LangChain 的Tool抽象和 Vercel AI SDK 的toolCall。bailian_web_search搜索服务对标 Google Programmable Search Engine 和 Bing Web Search API。
本节 AI 替代率:~85% | 人工干预率:~15%
| 角色 | 能力范围 |
|---|---|
| 🤖 AI 擅长 | MCP 客户端创建、SSE 连接代码、搜索调用、JSON 结果解析、异步桥接(asyncio.run)、容错处理 |
| 👤 人类需理解 | MCP 协议的工作原理(SSE + 工具调用 vs REST API)、异步编程的桥接方式(同步 process 中调用异步 _mcp_call) |
阅读指引
| 颜色 | 章节 | AI 替代率 | 人工干预 | 说明 |
|---|---|---|---|---|
| 🟡 | §1 任务目标 | ~95% | ~5% | 学习 MCP / SSE / async |
| 🟢 | §2 核心概念 | ~90% | ~10% | MCP 协议 / SSE / async 基础 |
| 🟡 | §3 整体流程 | ~90% | ~10% | 节点位置 + 输入输出 |
| 🟠 | §4 分步实现 | ~80% | ~20% | 6 步流程,Step3(MCP连接) + Step4(工具调用) 是核心 |
| 🔴 | §4.4 主代码 | ~75% | ~25% | WebSearchMcpNode + _mcp_call 异步桥接 |
| 🟢 | §5 测试运行 | ~95% | ~5% | 看预期输出即可 |
| 🟡 | §6 总结 | ~90% | ~10% | 设计要点回顾 |
技术栈健康度标签体系
| 技术 | 健康度 | 建议 |
|---|---|---|
| MCP 协议 | ⏳ 成长期 | Anthropic 2024 年底提出的 AI-工具通信协议。SSE 模式稳定,但生态仍在发展。2025-2026 年与 A2A(Google)竞争。 |
| MCPServerSse | ⏳ 成长期 | OpenAI Agents SDK 的 MCP 客户端实现。API 仍在迭代中。 |
| asyncio | 🔥 巅峰 | Python 异步框架标准库。asyncio.run() 桥接同步/异步是 Python 开发的通用模式。 |
| SSE (Server-Sent Events) | 🟢 稳定 | HTTP 单向推送技术,天然适配 AI 流式输出场景。自动重连机制成熟。 |
| Web Search (百炼) | 🔥 巅峰 | 阿里云百炼的 AI 搜索服务,MCP 协议接入。与 Bing / Google 搜索 API 等效。 |
体系说明:🟢🟡🟠🔴 标识学习优先级 / AI 替代率;🔥🟢⏳⚠️💀 标识技术栈健康度。
中英文对照表
| English | 中文 | 本质 |
|---|---|---|
| MCP (Model Context Protocol) | 模型上下文协议 | AI 模型调用外部工具的标准化通信协议 |
| SSE (Server-Sent Events) | 服务器推送事件 | HTTP 协议上的单向推送技术,服务端主动向客户端发送数据 |
| Tool Calling | 工具调用 | AI 模型选择并调用外部工具的机制 |
| asyncio.run() | 异步运行器 | 在同步代码中执行异步函数的桥接工具 |
| MCPServerSse | MCP SSE 服务器 | OpenAI Agents SDK 的 MCP 客户端实现,基于 SSE 连接 |
| Cleanup | 资源清理 | 释放 MCP 连接的 finally 保障机制 |
| Snippet | 摘要 | 搜索结果的文本片段摘要 |
| Timeout | 超时 | 连接/读取操作的最大等待时间(秒) |
💡 程序员比喻
- WebSearchMcpNode 就像
curl+jq——MCPServerSse= curl(发起 HTTP 请求),json.loads(result.content[0].text)= jq(解析 JSON)。- MCP 协议 就像 CI/CD 的 Trigger Token——一个标准化的"触发接口",我给你 token(API Key),你执行我的工具(搜索),返回结果给我。
- SSE 连接 就像
kubectl logs -f——服务端持续推送数据流,客户端收到后自动解析。- asyncio.run() 桥接 就像
subprocess.run('python script.py')——在 Java(同步)代码中调用 Python(异步)脚本。- finally + cleanup 就像
defer f.Close()in Go——不管函数正常返回还是 panic,资源都要释放。
1. 任务目标
本节课将实现知识库查询流程中的 网络搜索节点(web_search_mcp)。
该节点通过 MCP 协议(Model Context Protocol)调用外部网络搜索服务,获取实时的互联网搜索结果,作为本地知识库的补充信息源。
学完本节你将掌握:
- MCP 协议的基本概念与工作原理
- SSE(Server-Sent Events)连接模式
- Python 异步编程(async/await)在节点中的应用
- 如何将外部搜索服务集成到 LangGraph 工作流
2. 核心概念扫盲
2.1 为什么需要网络搜索?
本地知识库有以下局限性:
| 局限性 | 说明 |
|---|---|
| 时效性 | 知识库内容可能过时,无法回答最新问题 |
| 覆盖面 | 知识库只包含已导入的文档,无法覆盖所有领域 |
| 边界问题 | 用户可能提出超出知识库范围的问题 |
网络搜索作为"兜底"机制,可以:
- 补充知识库未覆盖的内容
- 提供最新的时效性信息
- 增加答案的多样性和可信度
2.2 MCP 协议是什么?
MCP(Model Context Protocol) 是一种标准化协议,用于 AI 模型与外部工具/服务之间的通信。

MCP 的核心优势:
- 标准化:统一的工具调用接口
- 解耦:AI 应用与工具实现分离
- 可扩展:轻松添加新工具而不修改 AI 应用
2.3 SSE 连接模式
SSE(Server-Sent Events) 是一种服务器向客户端推送事件的技术。

为什么使用 SSE?
- 单向通信,服务器主动推送
- 基于 HTTP,穿透防火墙能力强
- 自动重连,连接稳定
- 适合流式返回结果的场景
2.4 Python 异步编程基础
网络搜索涉及 I/O 等待,使用异步编程可以提高效率:
🟡 【P1 看注释就行】 Step1 获取查询——
rewritten_query,空则返回空 dict。
# 同步方式(阻塞)
def sync_search(query):
result = call_api(query) # 等待期间线程被阻塞
return result
# 异步方式(非阻塞)
async def async_search(query):
result = await call_api(query) # 等待期间可执行其他任务
return result
# 在同步代码中调用异步函数
result = asyncio.run(async_search(query))关键概念:
async def:定义异步函数(协程)await:等待异步操作完成asyncio.run():在同步环境中运行异步代码
3. 网络搜索业务处理流程(总)
3.1 节点在流程中的位置

web_search_mcp 与向量检索、知识图谱检索并行执行,共同为 RRF 融合提供候选结果。
3.2 节点输入输出

4. 网络搜索业务处理流程(分)
4.1 目标
实现一个通过 MCP 协议调用外部搜索服务的节点,获取与用户查询相关的网络搜索结果。
4.2 需求分析
| 需求项 | 说明 |
|---|---|
| MCP 连接 | 使用 SSE 模式连接到搜索服务 |
| 认证授权 | 通过 Header 传递 API Key |
| 超时控制 | 设置合理的连接和读取超时 |
| 结果解析 | 从 JSON 响应中提取有效字段 |
| 错误处理 | 网络异常时不影响整体流程 |
| 资源释放 | 确保连接正确关闭 |
4.3 实现流程
4.3.1 实现流程图

4.3.2 具体实现步骤
Step 1: 获取查询内容
从状态中获取重写后的查询文本:
🟡 【P1 看注释就行】 Step2 创建 MCP 客户端——
MCPServerSse的 4 个参数:url / headers / timeout / sse_read_timeout。注意timeout=300是 5 分钟超时。
query = state.get("rewritten_query", "")
if not query:
self.logger.warning("查询内容为空,跳过网络搜索")
return {}要点:
- 使用
rewritten_query而非original_query,因为重写后的查询更适合搜索 - 空查询直接返回空字典,不影响后续流程
Step 2: 创建 MCP 客户端
配置 SSE 连接参数并创建客户端实例:
🔥 【P0 必须要学】 Step3 SSE 连接——
await search_mcp.connect()建立长连接。MCP 的 SSE 模式与传统 REST API 的区别:连接后保持长连接,后续的 tool call 通过此连接传输,不是每次都重新建立 HTTP 连接。
from agents.mcp import MCPServerSse
dashscope_base_url = os.environ.get("MCP_DASHSCOPE_BASE_URL")
api_key = os.environ.get("OPENAI_API_KEY")
search_mcp = MCPServerSse(
name="search_mcp",
params={
"url": dashscope_base_url, # MCP 服务地址
"headers": {"Authorization": api_key}, # 认证头
"timeout": 300, # 连接超时(秒)
"sse_read_timeout": 300 # SSE 读取超时(秒)
}
)配置说明:
| 参数 | 说明 |
|---|---|
url | MCP 服务的 SSE 端点地址 |
headers | HTTP 请求头,用于传递认证信息 |
timeout | 建立连接的最大等待时间 |
sse_read_timeout | 读取 SSE 事件的最大等待时间 |
Step 3: 连接 MCP 服务
使用异步方式建立连接:
🔥 【P0 必须要学】 Step4 工具调用——
await search_mcp.call_tool(tool_name="bailian_web_search")。MCP 的标准调用模式:call_tool(tool_name, arguments)。注意count=5限制返回结果数。
await search_mcp.connect()连接过程:
- 向 MCP 服务发起 HTTP 请求
- 服务端返回 SSE 事件流
- 客户端保持长连接,准备接收事件
Step 4: 调用搜索工具
通过 MCP 协议调用搜索工具:
🟡 【P1 看注释就行】 Step5 解析搜索结果——
json.loads(result.content[0].text)获取 JSON,pages数组提取 title/url/snippet。过滤not snippet的无内容项。
result = await search_mcp.call_tool(
tool_name="bailian_web_search", # 工具名称
arguments={
"query": query, # 搜索关键词
"count": 5 # 返回结果数量
}
)MCP 工具调用流程:

Step 5: 解析搜索结果
从 MCP 响应中提取搜索结果:
🟡 【P1 看注释就行】 Step6 资源清理 + 返回——
finally: await search_mcp.cleanup()确保连接释放。失败时返回{}(空 dict),不影响其他字段。
# 1. 解析 JSON 响应
pages = json.loads(result.content[0].text).get("pages") or []
# 2. 提取有效字段
docs = []
for item in pages:
snippet = (item.get("snippet") or "").strip()
url = (item.get("url") or "").strip()
title = (item.get("title") or "").strip()
# 3. 过滤无效结果(必须有摘要)
if not snippet:
continue
docs.append({
"title": title,
"url": url,
"snippet": snippet
})响应结构示例:
{
"pages": [
{
"title": "数字万用表使用指南",
"url": "https://example.com/guide",
"snippet": "数字万用表测量电压时,首先将旋钮转到V档位..."
},
{
"title": "万用表入门教程",
"url": "https://example.com/tutorial",
"snippet": "测量直流电压的步骤如下:1. 选择合适的量程..."
}
]
}Step 6: 清理资源并返回
确保 MCP 连接正确关闭:
🔥 【P0 必须要学】 WebSearchMcpNode 完整代码——四个关键设计:(1)
asyncio.run(self._mcp_call(query))同步中调异步 (2)try/finally资源清理模式 (3) 空查询/失败返回{}的容错 (4)snippet空过滤。注意_mcp_call是 async 方法,process是同步方法——这个模式在其他异步 I/O 场景中可以复用。
try:
await search_mcp.connect()
result = await search_mcp.call_tool(...)
return result
finally:
# 无论成功失败,都要关闭连接
await search_mcp.cleanup()返回搜索结果到状态:
🟢 【P2 后面可以查】 测试代码——确保
MCP_DASHSCOPE_BASE_URL和OPENAI_API_KEY已配置。看预期输出即可。
if docs:
return {"web_search_docs": docs}
return {}4.4 代码实现
完整的节点实现代码:
"""MCP 网络搜索节点
通过 MCP 协议调用网络搜索服务获取外部信息。
"""
import asyncio
import os
import json
from knowledge.processor.query_process.base import BaseNode, setup_logging
from knowledge.processor.query_process.state import QueryGraphState
class WebSearchMcpNode(BaseNode):
"""MCP 网络搜索节点。
通过 MCP 协议连接到网络搜索服务,
根据用户查询获取相关的网络搜索结果。
"""
name = "web_search_mcp"
def process(self, state: QueryGraphState) -> QueryGraphState:
"""执行网络搜索。
Args:
state: 查询图状态,需包含 rewritten_query。
Returns:
更新后的状态,包含 web_search_docs 搜索结果。
"""
# Step 1: 获取查询内容
self.log_step("step_1", "获取查询内容")
query = state.get("rewritten_query", "")
docs = []
if not query:
self.logger.warning("查询内容为空,跳过网络搜索")
return {}
# Step 2-5: 执行 MCP 搜索
self.log_step("step_2", f"执行 MCP 搜索: {query}")
try:
result = asyncio.run(self._mcp_call(query))
if result:
pages = json.loads(result.content[0].text).get("pages") or []
for item in pages:
snippet = (item.get("snippet") or "").strip()
url = (item.get("url") or "").strip()
title = (item.get("title") or "").strip()
if not snippet:
continue
docs.append({
"title": title,
"url": url,
"snippet": snippet
})
self.log_step("step_3", f"搜索完成,返回 {len(docs)} 条结果")
except Exception as e:
self.logger.error(f"MCP 搜索失败: {e}")
# Step 6: 返回结果
if docs:
return {"web_search_docs": docs}
return {}
async def _mcp_call(self, query: str):
"""调用 MCP 搜索服务。
Args:
query: 搜索查询。
Returns:
MCP 搜索结果。
"""
from agents.mcp import MCPServerSse
# Step 2: 创建 MCP 客户端
dashscope_base_url = os.environ.get("MCP_DASHSCOPE_BASE_URL")
api_key = os.environ.get("OPENAI_API_KEY")
search_mcp = MCPServerSse(
name="search_mcp",
params={
"url": dashscope_base_url,
"headers": {"Authorization": api_key},
"timeout": 300,
"sse_read_timeout": 300
}
)
try:
# Step 3: 连接 MCP 服务
await search_mcp.connect()
# Step 4: 调用搜索工具
result = await search_mcp.call_tool(
tool_name="bailian_web_search",
arguments={"query": query, "count": 5}
)
return result
finally:
# Step 6: 清理资源
await search_mcp.cleanup()
# 兼容原有调用方式
_node_instance = WebSearchMcpNode()
def node_web_search_mcp(state: QueryGraphState) -> QueryGraphState:
"""MCP 网络搜索节点入口函数(兼容原有调用方式)。"""
return _node_instance(state)5. 测试运行
5.1 运行网络搜索节点测试
# 确保配置了必要的环境变量
export MCP_DASHSCOPE_BASE_URL="https://your-mcp-service-url"
export OPENAI_API_KEY="your-api-key"
# 运行测试
python -m knowledge.processor.query_process.nodes.web_search_mcp5.2 测试代码
if __name__ == "__main__":
from dotenv import load_dotenv
import json
# 1. 加载环境变量
load_dotenv()
# 2. 初始化日志
try:
setup_logging()
except Exception:
import logging
logging.basicConfig(level=logging.INFO)
print("=" * 60)
print("开始测试: MCP 网络搜索节点 (WebSearchMcpNode)")
print("=" * 60)
# 3. 构造模拟的初始状态
test_query = "数字万用表怎么测电压?"
initial_state = {
"original_query": test_query,
"rewritten_query": test_query
}
print("【输入状态】:")
print(json.dumps(initial_state, ensure_ascii=False, indent=2))
print("-" * 60)
try:
# 4. 执行网络搜索节点
final_state = node_web_search_mcp(initial_state)
# 5. 打印搜索结果
docs = final_state.get("web_search_docs", [])
if not docs:
print("\n搜索执行完成,但未返回任何结果。")
else:
print(f"\n执行完成!共搜索到 {len(docs)} 条相关内容:\n")
for i, doc in enumerate(docs, 1):
title = doc.get('title', '无标题')
url = doc.get('url', '无链接')
snippet = doc.get('snippet', '')
print(f"[{i}] 标题: {title}")
print(f" 链接: {url}")
print(f" 摘要: {snippet[:100]}...")
print("-" * 60)
except Exception as e:
print(f"\n执行失败: {e}")
import traceback
traceback.print_exc()5.3 预期输出
============================================================
开始测试: MCP 网络搜索节点 (WebSearchMcpNode)
============================================================
【输入状态】:
{
"original_query": "数字万用表怎么测电压?",
"rewritten_query": "数字万用表怎么测电压?"
}
------------------------------------------------------------
[web_search_mcp] [step_1] 获取查询内容
[web_search_mcp] [step_2] 执行 MCP 搜索: 数字万用表怎么测电压?
[web_search_mcp] [step_3] 搜索完成,返回 5 条结果
执行完成!共搜索到 5 条相关内容:
[1] 标题: 数字万用表使用方法详解
链接: https://example.com/guide/multimeter
摘要: 数字万用表测量电压的步骤:1. 将黑色表笔插入COM孔,红色表笔插入V孔;2. 将功能旋钮转到V...
------------------------------------------------------------
[2] 标题: 万用表测电压入门教程
链接: https://example.com/tutorial
摘要: 测量直流电压时,将旋钮转到DCV档位,选择合适的量程,然后将表笔并联到被测电路两端...
------------------------------------------------------------
...5.4 处理前后对比
| 对比项 | 处理前 | 处理后 |
|---|---|---|
rewritten_query | "数字万用表怎么测电压?" | 不变 |
web_search_docs | [] 或不存在 | 包含 5 条搜索结果 |
| 每条结果结构 | - | {title, url, snippet} |
数据结构变化:
# 处理前
state = {
"rewritten_query": "数字万用表怎么测电压?",
# web_search_docs 不存在
}
# 处理后
state = {
"rewritten_query": "数字万用表怎么测电压?",
"web_search_docs": [
{
"title": "数字万用表使用方法详解",
"url": "https://example.com/guide",
"snippet": "数字万用表测量电压的步骤..."
},
{
"title": "万用表测电压入门教程",
"url": "https://example.com/tutorial",
"snippet": "测量直流电压时,将旋钮转到..."
},
# ... 更多结果
]
}6. 总结
6.1 节点功能概览
┌─────────────────────────────────────────────────────────────┐
│ WebSearchMcpNode │
├─────────────────────────────────────────────────────────────┤
│ │
│ 核心功能: 通过 MCP 协议调用网络搜索服务 │
│ │
│ 输入: │
│ └── state["rewritten_query"] 重写后的查询 │
│ │
│ 输出: │
│ └── state["web_search_docs"] 网络搜索结果列表 │
│ ├── title 网页标题 │
│ ├── url 网页链接 │
│ └── snippet 内容摘要 │
│ │
│ 依赖: │
│ ├── agents.mcp.MCPServerSse MCP 客户端 │
│ ├── MCP_DASHSCOPE_BASE_URL 服务地址 │
│ └── OPENAI_API_KEY 认证密钥 │
│ │
│ 特点: │
│ ├── 异步执行,不阻塞主流程 │
│ ├── 容错设计,失败时返回空结果 │
│ └── 资源自动清理(finally) │
│ │
└─────────────────────────────────────────────────────────────┘6.2 节点设计要点
要点 1:MCP 协议的使用
# MCP 客户端配置
search_mcp = MCPServerSse(
name="search_mcp",
params={
"url": dashscope_base_url, # 服务端点
"headers": {"Authorization": api_key}, # 认证信息
"timeout": 300, # 连接超时
"sse_read_timeout": 300 # 读取超时
}
)
# 调用工具
result = await search_mcp.call_tool(
tool_name="bailian_web_search",
arguments={"query": query, "count": 5}
)要点 2:同步节点中调用异步代码
def process(self, state):
# 同步方法
...
# 使用 asyncio.run() 桥接异步调用
result = asyncio.run(self._mcp_call(query))
...
async def _mcp_call(self, query):
# 异步方法
await search_mcp.connect()
result = await search_mcp.call_tool(...)
return result要点 3:资源清理模式
try:
await search_mcp.connect()
result = await search_mcp.call_tool(...)
return result
finally:
# 确保连接被关闭,避免资源泄漏
await search_mcp.cleanup()要点 4:容错设计
try:
result = asyncio.run(self._mcp_call(query))
# 处理结果...
except Exception as e:
# 记录错误但不抛出,避免影响整体流程
self.logger.error(f"MCP 搜索失败: {e}")
# 无论成功失败,都返回有效状态
if docs:
return {"web_search_docs": docs}
return {} # 空字典不会影响其他字段要点 5:结果过滤
for item in pages:
snippet = (item.get("snippet") or "").strip()
# ...
# 只保留有有效摘要的结果
if not snippet:
continue
docs.append({...})企业痛点映射
| 痛点 | 传统方案 | MCP 网络搜索方案 | 效率提升 |
|---|---|---|---|
| 知识库内容过时 | 手动更新文档,周期长 | 实时网络搜索补充最新信息 | 时效性信息覆盖从 ~0% 到 实时 |
| 知识库未覆盖的内容 | 返回"抱歉,无法回答" | MCP 调用搜索引擎兜底 | 回答覆盖率提升 ~30%(预估) |
| 异步代码与同步流程不兼容 | 无法在 LangGraph 同步节点中调用异步 API | asyncio.run() 桥接 | 零代码改造,兼容 LangGraph |
| MCP 连接泄漏导致资源耗尽 | 每次搜索都创建新连接不关闭 | try/finally + cleanup() 保障资源释放 | 连接泄漏风险降为 ~0% |
Remote & Agent 应用场景价值
Remote 场景价值:MCP 协议天然适配远程架构——搜索服务在远端部署,本地通过 SSE 连接。
timeout=300的超时配置适合不稳定的远程网络环境。异步非阻塞设计让团队成员在 CI/CD 中执行搜索时不阻塞其他测试用例。Agent 落地场景:WebSearchMcpNode 可作为"网络搜索 Agent"——Agent 收到查询 → MCP 协议调用搜索引擎 → 解析结果 → 返回结构化文档。在 Multi-Agent 架构中,该 Agent 是"信息检索团队"的成员,与向量检索 Agent、KG Agent 并行工作,结果由 RRF Agent 融合。MCP 协议的标准化接口使 Agent 可以无缝切换搜索引擎(百炼 → Bing → Google),只需改
tool_name。
Git Commit 对应
本节网络搜索节点对应的提交记录(参考值,以实际版本为准):
<待补充 — 建议搜索 "web_search_mcp.py" 相关提交>cd shopkeeper_brain
git log --oneline --all -- knowledge/processor/query_process/nodes/web_search_mcp.py