之前在做 AI Agent 落地时,我一直被一个看似基础的问题困扰:Agent 明明可以调用多个工具,但每次让它“查资料”时,效果总是不稳定。要么是单个工具返回内容太少,要么是多个工具的结果重复且格式混乱。直到我把搜索能力从 Agent 内部拆出来,单独做成一个联邦搜索服务后,整条链路才真正顺畅起来。本文就来完整拆解 Federated Search for AI Agents 的架构思路、代码实现和工程落地经验。
不论你是在做 RAG 应用、智能客服,还是让 Agent 对接内部多个业务系统的数据,这套方案都能直接复用。接下来我们会先明确联邦搜索到底解决了什么问题,再逐步搭建一个带完整代码的检索服务,最后把它封装成 Agent 可以调用的工具。
1. 为什么 AI Agent 需要一个“联邦搜索”层
1.1 Agent 的检索难题:工具越多,信息越碎
很多同学在初学 Agent 时,会先做一个“超级工具人”:给 Agent 挂上搜索引擎、数据库查询、内部 API、向量库检索等一大堆工具。看起来能力很全,但真正跑起来会遇到几个扎心的问题。
第一个问题是信息孤岛。每个工具只负责自己的数据源,而用户的问题通常需要跨多个源才能回答。比如“帮我查一下最近一周线上问题中,哪些和支付超时有关,同时附上对应的工单处理状态”。这个问题同时涉及日志系统、工单系统、知识库,如果 Agent 按顺序逐个调用工具,很容易丢上下文,也可能把不同源的信息拼接得前言不搭后语。
第二个问题是结果爆炸。多个工具各自返回内容,加起来的 token 数量可能远超模型上下文窗口。如果直接把大量原始结果塞给 LLM,不仅浪费 token,还会稀释关键信息,导致 Agent 回答偏题。
第三个问题是质量不稳定。不同数据源的返回格式不统一,有的是 JSON,有的是 HTML,有的是纯文本,还有的可能返回空结果或超时。Agent 自己处理这些异常时,Prompt 稍微设计得不够健壮,就会出现“工具调用报错”或者“答非所问”。
1.2 从查询到答案:联邦搜索在这条链路中的位置
联邦搜索(Federated Search)并不是一个新概念,它在传统信息检索领域已经存在很久了。简单来说,它指的是把同一个查询请求同时分发到多个独立的检索源,再把各源返回的结果进行统一整合,最终给用户一个合并后的结果列表。
在 AI Agent 场景下,联邦搜索的角色有点像一个“搜索中台”:
用户问题 -> Agent 规划 -> 联邦搜索服务 -> 数据源A -> 数据源B -> 数据源C <- 合并去重排序 Agent 拿到结构化结果 -> LLM 生成答案它把“查哪里”和“怎么查”的复杂度从 Agent 内部抽离出来,让 Agent 只需要关心一个统一的检索入口。这样做的好处很明显:
- Agent 不需要在 Prompt 里罗列几十个工具。
- 每个数据源可以独立维护,新增数据源不影响 Agent。
- 结果经过合并、去重、排序,质量更高,token 消耗更少。
- 超时、降级、鉴权等逻辑统一收敛到服务层。
1.3 联邦搜索不是“调用多个接口”这么简单
有些人可能会说:不就是用 Python 写个循环,然后逐个调 API 吗?
其实联邦搜索的核心难点不在“并发调用”,而在资源描述和结果融合。
资源描述解决的是“查哪些源”的问题。比如有 10 个搜索源,用户问天气时显然不需要查工单系统,用户问 Bug 时大概率不需要查天气 API。如果每次查询都全量分发,既浪费资源,又让结果噪音变大。
结果融合解决的是“结果怎么排”的问题。不同数据源的评分体系不一样,A 源返回 0-100 分,B 源返回 0-1 的相似度,C 源干脆不返回分数。如果不做归一化和加权,直接按原始顺序拼接,排序就是乱的。
所以我们构建的联邦搜索服务,至少要包含四个部分:
- 查询理解与源选择。
- 并发请求调度。
- 结果归一化与融合。
- 统一响应输出。
2. 联邦搜索的几个核心概念
2.1 Query Dispatcher:查询分发
查询分发器负责决定把请求发给哪些源,以及如何并发执行。
常见的策略有三种:
- 全量分发:所有源都查,适合源数量少、查询成本低的场景。
- 规则分发:根据关键词或业务属性,通过配置决定命中哪些源。
- 智能选择:借助 LLM 或轻量分类模型,让模型从源列表中挑选相关源。
在工程上,我建议先做“规则分发 + 全量兜底”。等数据源多了,再引入 LLM 做源选择,因为 LLM 选择源这件事本身也有延迟和误判成本。
2.2 Resource Description:资源描述与选择
每个搜索源都需要有一份元数据描述,告诉联邦搜索服务:
- 这个源是什么?(名称、类型)
- 它适合处理什么类型的问题?(领域、标签)
- 它返回什么结构?(字段映射)
- 它的性能和稳定性如何?(超时时间、优先级)
在代码中,我们可以用 dataclass 或 Pydantic 模型来表达资源描述,后续做源选择时,本质上就是遍历这些描述做匹配。
2.3 Result Fusion:结果融合
结果融合是联邦搜索中最经典的研究方向。常见的融合算法有:
- 加权评分融合:对每个源的结果评分做归一化,再乘以权重,最终按加权分数排序。
- Borda Count:每个源选出 Top-K 结果,根据排名位置计分,按总得分排序。
- Reciprocal Rank Fusion:每个结果按其在不同源中的排名计算
1/(k + rank),累加得分。
对 Agent 场景来说,推荐的方案是“加权评分 + 归一化”,因为它实现简单,而且权重可解释、可控。如果某个源没有分数,可以用排序位置倒推出一个分数。
2.4 与 RAG、Agent 工具调用的关系
很多人分不清联邦搜索和 RAG 的区别。
RAG(检索增强生成)通常指的是:把一个查询拿去检索自己的知识库(通常是向量库),把检索到的文本片段塞进 Prompt,帮助 LLM 回答。RAG 的检索源相对单一,聚焦在“文档召回”。
联邦搜索则是面向多个异构数据源的统一检索入口。它更强调“跨源查询”和“结果合并”。
两者并不冲突,可以组合使用。比如一个联邦搜索服务,内部可以同时接入一个向量库、一个 SQL 数据库、一个外部 API。向量库负责语义召回,数据库负责结构化查询,API 负责实时数据。这样一来,Agent 拿到的结果既有文档片段,又有结构化数据,信息完整度会高很多。
3. 环境准备与项目结构
3.1 运行环境与依赖
本文的实战代码以 Python 3.10+ 为例,核心依赖如下:
fastapi:提供检索服务的 HTTP 接口。uvicorn:ASGI 服务器,用于启动 FastAPI。httpx:异步 HTTP 客户端,用于并发调用外部 API。pydantic:数据模型定义与校验。langchain(可选):用来把联邦搜索服务封装成 Agent 工具。
版本不需要刻意选择最新版,以你本机环境兼容为准。如果不想引入 LangChain,只要 Agent 框架支持 HTTP 调用,也可以直接把 FastAPI 接口挂进去。
建议使用虚拟环境安装:
python -m venv venv source venv/bin/activate # Windows 下使用 venv\Scripts\activate pip install fastapi uvicorn httpx pydantic如果你希望同时跑 LangChain 示例,再安装:
pip install langchain langchain-openai3.2 示例项目结构
为了便于阅读,我们把代码拆成多个模块。下文示例中会使用如下项目结构:
federated_search/ ├── main.py # FastAPI 入口 ├── schemas.py # 请求与响应模型 ├── sources.py # 数据源适配器 ├── fusion.py # 结果融合逻辑 ├── dispatcher.py # 查询分发器 └── agent_tool.py # 封装成 Agent 工具4. 实战:搭建一个面向 Agent 的联邦搜索服务
下面我们逐步实现一个最小可运行的联邦搜索服务。为了演示,我会模拟两个数据源:一个“技术文档知识库”,一个“工单系统 API”。实际项目中,你可以把这两个适配器替换成真实的接口。
4.1 定义统一的检索请求与响应模型
不同数据源返回格式五花八门,所以我们首先要定义一套中间格式,让所有源适配器都输出这个格式。
# 文件路径:federated_search/schemas.py from typing import Optional from pydantic import BaseModel, Field class SearchQuery(BaseModel): """联邦搜索的统一查询请求""" query: str = Field(..., description="用户查询内容") top_k: int = Field(10, ge=1, le=50, description="每个数据源返回的最大条数") sources: Optional[list[str]] = Field(None, description="指定要查询的数据源,默认查全部") class SearchResultItem(BaseModel): """单个检索结果项""" source: str = Field(..., description="数据源名称") title: str = Field("", description="结果标题") content: str = Field("", description="结果内容或摘要") url: str = Field("", description="跳转链接或唯一标识") score: float = Field(0.0, description="原始分数") raw: dict = Field(default_factory=dict, description="原始数据,方便调试") class SearchResponse(BaseModel): """联邦搜索的统一响应""" query: str results: list[SearchResultItem] trace: dict = Field(default_factory=dict, description="各源调用情况,方便排查")这里使用 Pydantic 模型,既能在接口层做参数校验,又能保证结果结构统一。
4.2 实现多个数据源适配器
每个数据源都实现一个SearchSource接口。这里我们定义一个基类,然后写两个具体适配器。
第一个模拟技术文档知识库,返回预设的结果。第二个模拟工单系统 API,通过 httpx 调用真实接口(在示例中我直接模拟,因为真实接口不可用)。
# 文件路径:federated_search/sources.py from abc import ABC, abstractmethod import asyncio import httpx from schemas import SearchQuery, SearchResultItem class BaseSource(ABC): """所有数据源适配器的基类""" name: str = "base" def __init__(self, timeout: float = 5.0): self.timeout = timeout @abstractmethod async def search(self, query: SearchQuery) -> list[SearchResultItem]: """执行搜索,返回统一格式的结果列表""" raise NotImplementedError class DocKnowledgeSource(BaseSource): """模拟技术文档知识库""" name = "docs" def __init__(self): super().__init__(timeout=2.0) async def search(self, query: SearchQuery) -> list[SearchResultItem]: # 模拟网络 IO 延迟 await asyncio.sleep(0.2) keywords = query.query # 这里应该是向量检索或数据库查询,示例中直接构造模拟数据 mock_results = [ SearchResultItem( source=self.name, title="FastAPI 官方文档:请求体", content=f"FastAPI 中定义请求体可以使用 Pydantic 模型,关键词:{keywords}", url="https://example.com/fastapi-body", score=0.92, ), SearchResultItem( source=self.name, title="httpx 异步客户端使用指南", content="httpx 支持 async / await,可以并发发起多个请求,关键词:" + keywords, url="https://example.com/httpx-async", score=0.87, ), ] return mock_results[: query.top_k] class TicketSystemSource(BaseSource): """模拟工单系统 API""" name = "ticket" def __init__(self, api_base: str = ""): super().__init__(timeout=3.0) self.api_base = api_base async def search(self, query: SearchQuery) -> list[SearchResultItem]: # 如果配了真实 api_base,就用 httpx 调用 if self.api_base: async with httpx.AsyncClient(timeout=self.timeout) as client: resp = await client.post( f"{self.api_base}/search", json={"query": query.query, "top_k": query.top_k}, ) data = resp.json() return [ SearchResultItem( source=self.name, title=item.get("title", ""), content=item.get("content", ""), url=item.get("url", ""), score=float(item.get("score", 0.0)), raw=item, ) for item in data.get("items", []) ] # 模拟返回 await asyncio.sleep(0.3) mock_results = [ SearchResultItem( source=self.name, title=f"工单 #1024:登录超时问题", content="用户在登录时遇到超时,与查询相关:" + query.query, url="https://ticket.example.com/1024", score=0.88, ), SearchResultItem( source=self.name, title=f"工单 #2048:支付回调异常", content="支付回调处理失败,可能需要联动排查:" + query.query, url="https://ticket.example.com/2048", score=0.79, ), ] return mock_results[: query.top_k]这里有一个值得注意的点:SearchResultItem中我特意加入了raw字段,因为它对排查问题非常有用。生产环境里,各源返回的原始数据可能和标准字段对不上,保留raw可以避免丢失信息。
4.3 实现查询分发与并发调用
查询分发器的职责是:确定要调用哪些源,并发执行搜索,收集结果,最后调用融合模块。
# 文件路径:federated_search/dispatcher.py import asyncio from schemas import SearchQuery, SearchResponse from sources import BaseSource class Dispatcher: def __init__(self, sources: dict[str, BaseSource]): self.sources = sources def _select_sources(self, query: SearchQuery) -> list[BaseSource]: """选择数据源:优先使用 query.sources,否则全量分发""" if query.sources: selected = [] for name in query.sources: if name in self.sources: selected.append(self.sources[name]) else: print(f"[warn] 未知数据源: {name}") return selected return list(self.sources.values()) async def dispatch(self, query: SearchQuery) -> SearchResponse: selected_sources = self._select_sources(query) tasks = [source.search(query) for source in selected_sources] # 使用 gather 并发执行,return_exceptions=True 避免单个源异常影响整体 results_list = await asyncio.gather(*tasks, return_exceptions=True) all_results = [] trace = {} for source, results in zip(selected_sources, results_list): if isinstance(results, Exception): # 生产环境这里应该接入日志系统 print(f"[error] 数据源 {source.name} 查询失败: {results}") trace[source.name] = {"status": "error", "error": str(results)} continue trace[source.name] = {"status": "ok", "count": len(results)} all_results.extend(results) return SearchResponse(query=query.query, results=all_results, trace=trace)这段代码的关键是asyncio.gather(..., return_exceptions=True)。它保证了某个数据源超时或报错时,不会拖垮整个联邦搜索服务。实际项目中,这里还可以加入“熔断”机制,比如某个源连续失败 N 次后自动摘除。
4.4 实现结果融合与去重
结果融合是联邦搜索中最重要的算法环节。这里实现一个轻量级方案:
- 对每个结果按源内排名计算一个归一化分数。
- 如果源自带 score,也一并归一化到 0-1。
- 按“源权重 + 归一化分数”计算最终分。
- 按标题做简单去重,合并相同内容。
# 文件路径:federated_search/fusion.py from collections import defaultdict from schemas import SearchResultItem def normalize_score(item: SearchResultItem, index: int, total: int) -> float: """返回一个 0-1 的分数""" # 如果源自带可信分数,使用分数;否则用位置倒推 if item.score > 0: return max(0.0, min(1.0, item.score / 100.0)) if total <= 1: return 1.0 return max(0.0, 1.0 - index / (total - 1)) def title_key(title: str) -> str: """简单的标题规范化,用于去重""" return title.strip().lower().replace(" ", "") def fuse_results( results: list[SearchResultItem], top_k: int = 10, source_weights: dict[str, float] | None = None, ) -> list[SearchResultItem]: weights = source_weights or {} default_weight = 1.0 # 按源分组,计算每个源内部的位置排名分 grouped: dict[str, list[SearchResultItem]] = defaultdict(list) for item in results: grouped[item.source].append(item) # 统计每个源的总数,用于归一化 source_total = {source: len(items) for source, items in grouped.items()} merged: dict[str, SearchResultItem] = {} for source, items in grouped.items(): for index, item in enumerate(items): score = normalize_score(item, index, source_total[source]) weight = weights.get(source, default_weight) final_score = score * weight # 去重逻辑:标题相似则合并,取分数高者 key = title_key(item.title) if key in merged: merged[key].score = max(merged[key].score, final_score) # 保留内容更长、信息更全的结果 if len(item.content) > len(merged[key].content): merged[key] = item merged[key].score = final_score else: item.score = final_score merged[key] = item fused = list(merged.values()) fused.sort(key=lambda x: x.score, reverse=True) return fused[:top_k]这段代码在工程上比较实用。比如“文档知识库”的结果权威度更高,可以设置更高的权重;工单系统结果时效性更重要,可以调整权重比例。源权重最好放到配置中心或环境变量里,不要写死在代码中。
4.5 暴露 FastAPI 接口
把上述模块组装起来,写一个 FastAPI 入口。
# 文件路径:federated_search/main.py from fastapi import FastAPI from dispatcher import Dispatcher from fusion import fuse_results from schemas import SearchQuery, SearchResponse from sources import DocKnowledgeSource, TicketSystemSource app = FastAPI(title="Federated Search for AI Agents") # 初始化数据源 sources = { "docs": DocKnowledgeSource(), "ticket": TicketSystemSource(), } dispatcher = Dispatcher(sources) # 源权重,按业务需要调整 SOURCE_WEIGHTS = { "docs": 1.0, "ticket": 0.9, } @app.post("/search", response_model=SearchResponse) async def search_api(query: SearchQuery): response = await dispatcher.dispatch(query) response.results = fuse_results( response.results, top_k=query.top_k, source_weights=SOURCE_WEIGHTS, ) return response @app.get("/health") async def health(): return {"status": "ok", "sources": list(sources.keys())}启动服务:
cd federated_search uvicorn main:app --host 0.0.0.0 --port 8000然后用 curl 测试:
curl -X POST http://localhost:8000/search \ -H "Content-Type: application/json" \ -d '{"query": "支付超时问题", "top_k": 5}'预期会返回一个合并后的 JSON,结果已经按照融合分数排好序,并且两个模拟源的数据都出现在列表中。
5. 把联邦搜索封装成 Agent 的工具
联邦搜索服务本身已经是一个可用的 HTTP API,但真正要用到 Agent 中,还需要封装成 Agent 的工具。这里我以 LangChain 为例演示。
5.1 使用 LangChain 的 Tool 接口接入
LangChain 提供了一个tool装饰器,可以快速把一个函数变成可供 Agent 调用的工具。
# 文件路径:federated_search/agent_tool.py from langchain_core.tools import tool import httpx import json @tool def federated_search(query: str, top_k: int = 5) -> str: """ 在多个数据源中执行联邦搜索,返回合并排序后的结果。 适合查询技术文档、工单系统、内部知识库等场景。 """ response = httpx.post( "http://localhost:8000/search", json={"query": query, "top_k": top_k}, timeout=10.0, ) response.raise_for_status() data = response.json() lines = [] for idx, item in enumerate(data.get("results", []), start=1): lines.append( f"{idx}. [{item['source']}] {item['title']}\n" f" {item['content']}\n" f" {item['url']}" ) return "\n".join(lines) if __name__ == "__main__": # 本地测试工具返回 result = federated_search.invoke({"query": "支付超时", "top_k": 3}) print(result)这里需要注意,LangChain 不同版本对tool装饰器的参数支持有差异。以上代码基于较新的langchain_core.tools写法。如果遇到导入失败,可以按你当前版本的官方文档调整。
5.2 给 Agent 一个更“会用”的查询入口
将联邦搜索暴露给 Agent 后,还有几个细节需要注意。
- 工具描述要写得足够清楚。LLM 选择工具时,主要靠 description 判断。描述里要说明这个工具适合查什么、不适合查什么。
- 返回内容要结构化。工具返回给 LLM 的文本不能太长。建议只保留 title、source、url 和一段摘要。如果 Agent 需要更多细节,可以继续调用其他工具。
- 不要把所有结果都塞回去。按相关度截断,通常返回 5-8 条就够了。如果 Agent 需要更多,可以让它再调用一次。
5.3 一个完整的 Agent 调用示例
假设我们使用 LangChain 的标准 Agent 流程:
# 文件路径:federated_search/run_agent.py from langchain_openai import ChatOpenAI from langchain.agents import create_tool_calling_agent, AgentExecutor from langchain_core.prompts import ChatPromptTemplate from agent_tool import federated_search # 这里仅示例,请替换为你的模型配置 llm = ChatOpenAI(base_url="https://your-api.example.com/v1", api_key="your-key", model="gpt-4o-mini") tools = [federated_search] prompt = ChatPromptTemplate.from_messages([ ("system", "你是一个智能助手,可以调用联邦搜索工具获取跨源信息。"), ("human", "{input}"), ("placeholder", "{agent_scratchpad}"), ]) agent = create_tool_calling_agent(llm, tools, prompt) executor = AgentExecutor(agent=agent, tools=tools, verbose=True) response = executor.invoke({"input": "帮我查一下最近有没有支付超时相关的工单和技术文档?"}) print(response["output"])在这个示例中,Agent 根据用户问题判断需要查询多个数据源,然后调用联邦搜索工具,拿到合并后的结构化结果,再组织语言回答。整个过程,Agent 不需要关心联邦搜索服务内部有几个源、各自是什么协议,复杂度被完全隔离。
6. 常见问题与排查思路
在做联邦搜索时,大家最常遇到下面几个问题。我整理成了一张排查表:
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
| 单个数据源出错后整个搜索失败 | 没有对异常做隔离 | 在asyncio.gather中使用return_exceptions=True,并对单源异常做降级 |
| 结果排序混乱,相关度低 | 不同源分数不在同一量纲 | 先归一化分数,再按源权重加权融合 |
| 重复结果太多 | 多个源有相同内容 | 按标题规范化后进行去重,必要时用内容指纹去重 |
| 查询延迟过高 | 最慢的数据源拖慢整体响应 | 给每个源单独设置超时时间,并对慢源做缓存 |
| Agent 回答不准确 | 工具返回内容太长,关键信息被稀释 | 在工具内做摘要截断,只返回 Top-K 条结构化结果 |
| 新增数据源后结果变差 | 源选择策略没有及时调整 | 补充资源描述,调整源权重,必要时做离线评测 |
针对“某个源超时拖慢整体响应”的问题,可以按下面的方式处理:
- 给每个源设置独立的
timeout参数。 - 把查询请求放入
asyncio.wait_for,控制最大等待时间。 - 对响应慢的源增加缓存策略。比如文档源结果可以缓存 5 分钟,工单源结果缓存 1 分钟。
- 如果长期超时,考虑降级:从该源改为返回空结果,并记录错误日志。
另外一个容易被忽略的问题是接口鉴权。联邦搜索服务内部会访问多个数据源,这些源的 API Key、账号信息不能直接暴露给 Agent。最佳做法是把敏感信息放在服务端环境变量或配置中心,Agent 只访问联邦搜索服务一个入口,由服务层统一做鉴权和数据脱敏。
7. 最佳实践与工程化建议
7.1 配置管理
联邦搜索服务的数据源列表、超时时间、权重、缓存策略,都应该做成配置,而不是写死在代码里。推荐使用环境变量或配置中心管理。
一个便于维护的配置示例:
# config.yaml 或环境变量注入 sources: docs: type: http url: https://docs.example.com/search timeout: 2.0 weight: 1.0 cache_ttl: 300 ticket: type: http url: https://ticket.example.com/search timeout: 3.0 weight: 0.9 cache_ttl: 60这样每次调整某个源的权重或超时,不需要重新发布代码。
7.2 日志与链路追踪
联邦搜索是一个典型的多依赖服务,最适合用链路追踪来排查问题。建议记录以下信息:
- 请求 ID。
- 查询内容(注意脱敏)。
- 命中了哪些源。
- 每个源的响应耗时。
- 每个源返回的数量。
- 合并后的最终结果数量。
有了这些信息,当 Agent 回答质量不佳时,你可以快速判断是“源没查对”还是“融合排序不合理”。
7.3 缓存策略
联邦搜索中很多数据源是相对静态的,比如技术文档、知识库。对这些源做缓存可以显著降低延迟和外部 API 成本。
但要注意缓存一致性。比如工单系统的数据可能实时变化,缓存时间要短。文档源的内容更新频率低,缓存时间可以长一些。建议采用“按源配置 TTL”的方式,而不是统一缓存时间。
7.4 安全边界
让 Agent 调用多个数据源,本质上扩大了系统的攻击面。需要注意以下几点:
- 最小权限:联邦搜索服务用到的数据源账号,只授予它真正需要的读取权限。
- 数据脱敏:返回给 Agent 或用户的结果中,去掉不必要的敏感字段。
- 防止 Prompt 注入:如果查询词可能包含恶意指令,不要把查询词直接拼接进新的 Prompt,尽量让工具返回结构化数据。
- 结果源标识:在返回结果中保留 source 字段,方便审计和追溯。
7.5 可评测性
联邦搜索做得好不好,不能只靠感觉。建议建立一套评测集,包含典型问题、预期来源、预期排序。每次调整融合算法或源权重后,跑一遍评测集,对比排序质量。
评测指标可以用:
- P@K:前 K 条结果中相关结果的比例。
- MRR:第一个相关结果的排名倒数。
- nDCG:考虑排序位置的归一化折损累计增益。
这是联邦搜索和 Agent 结合后最容易忽略的部分。没有评测,权重调优就是在“猜”。
8. 后续可以深入的方向
到这里,我们已经实现了一个完整的“面向 AI Agent 的联邦搜索服务”。整体结构并不复杂,核心就是把异构数据源统一接入、并发分发、融合排序,最终变成 Agent 可以直接消费的干净结果。
如果你准备在真实项目中使用这套方案,我建议按这个顺序逐步推进:
先跑通最小闭环,也就是本文的 FastAPI 服务,加上 LangChain 工具封装。接着完善结果融合逻辑,引入源权重和去重策略。再接入真实的数据源,替换掉模拟适配器。最后再考虑智能源选择、缓存、链路追踪和评测体系。
如果后续数据源数量继续增加,可以继续学习这些方向:
- 基于 LLM 的源选择:让模型根据用户问题自动决定查哪些源。
- 个性化排序:根据用户画像或历史行为调整源权重。
- 联邦搜索与 RAG 深度融合:把向量检索结果也作为联邦搜索的一个数据源。
你可以把本文的工程实现看作一个可持续迭代的起点。先从最基础的两个源跑通,再逐步扩展。等联邦搜索层稳定后,你会发现 Agent 的回答质量和可维护性都会明显提升。