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 的 toolCallbailian_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_URLOPENAI_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 超级个体实战指南