Skip to content

知识库查询 —— 网络搜索节点 ​

本文档详细介绍知识库查询流程中的网络搜索节点(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()异步运行器在同步代码中执行异步函数的桥接工具
MCPServerSseMCP 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。

python
# 同步方式(阻塞)
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 分钟超时。

python
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 连接。

python
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 读取超时(秒)
    }
)

配置说明:

参数说明
urlMCP 服务的 SSE 端点地址
headersHTTP 请求头,用于传递认证信息
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 限制返回结果数。

python
await search_mcp.connect()

连接过程:

  1. 向 MCP 服务发起 HTTP 请求
  2. 服务端返回 SSE 事件流
  3. 客户端保持长连接,准备接收事件
Step 4: 调用搜索工具 ​

通过 MCP 协议调用搜索工具:

🟡 【P1 看注释就行】 Step5 解析搜索结果——json.loads(result.content[0].text) 获取 JSON,pages 数组提取 title/url/snippet。过滤 not snippet 的无内容项。

python
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),不影响其他字段。

python
# 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
    })

响应结构示例:

json
{
  "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 场景中可以复用。

python
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 已配置。看预期输出即可。

python
if docs:
    return {"web_search_docs": docs}
return {}

4.4 代码实现 ​

完整的节点实现代码:

python
"""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 运行网络搜索节点测试 ​

bash
# 确保配置了必要的环境变量
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_mcp

5.2 测试代码 ​

python
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}

数据结构变化:

python
# 处理前
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 协议的使用 ​

python
# 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:同步节点中调用异步代码 ​

python
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:资源清理模式 ​

python
try:
    await search_mcp.connect()
    result = await search_mcp.call_tool(...)
    return result
finally:
    # 确保连接被关闭,避免资源泄漏
    await search_mcp.cleanup()

要点 4:容错设计 ​

python
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:结果过滤 ​

python
for item in pages:
    snippet = (item.get("snippet") or "").strip()
    # ...

    # 只保留有有效摘要的结果
    if not snippet:
        continue

    docs.append({...})

企业痛点映射 ​

痛点传统方案MCP 网络搜索方案效率提升
知识库内容过时手动更新文档,周期长实时网络搜索补充最新信息时效性信息覆盖从 ~0% 到 实时
知识库未覆盖的内容返回"抱歉,无法回答"MCP 调用搜索引擎兜底回答覆盖率提升 ~30%(预估)
异步代码与同步流程不兼容无法在 LangGraph 同步节点中调用异步 APIasyncio.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" 相关提交>
bash
cd shopkeeper_brain
git log --oneline --all -- knowledge/processor/query_process/nodes/web_search_mcp.py

OPC 超级个体实战指南