很多做 Agent 开发的同学,一开始都是从单数据源检索入手的:一个向量库、一套文档、一个搜索 API。但当 Agent 真正要面对企业级场景时会发现,回答一个问题经常要同时查 Wiki、查数据库、查工单系统、查外部知识库,一个个工具串行调用,延迟高、结果杂、还容易遗漏。
这次我们来看一种面向 AI Agents 的检索架构——Federated Search for AI Agents,也就是把“联邦搜索”能力嵌入到 Agent 的工具链中。简单来说,它让 Agent 通过一个统一入口,同时查询多个异构数据源,再把结果合并、去重、重排后交给大模型生成回答。这篇文章会从架构拆解、最小实现、API 设计、批量任务、性能观察和排查思路几个维度展开,目标是让你对联邦搜索在 Agent 场景里的落地方式有一个完整判断。
1. 核心能力速览
| 能力项 | 说明 |
|---|---|
| 项目类型 | AI Agent 检索基础设施 / 多源搜索架构方案 |
| 核心目标 | 让 AI Agent 通过统一接口同时查询多个异构数据源 |
| 主要功能 | 多源查询规划、并发请求分发、结果去重与重排、统一返回格式 |
| 典型数据源 | SQL 数据库、向量数据库、Elasticsearch、Web API、文件索引等 |
| 与 RAG 的关系 | 可以理解为 RAG 检索层的“上游网关”,也能直接作为 Agent Tool 使用 |
| 是否支持批量任务 | 支持,批量查询需要额外设计任务队列和并发控制 |
| 是否提供 API | 取决于具体实现,通常以 HTTP/GRPC 服务形式暴露 |
| 推荐基础环境 | Linux 服务器或本地开发机,Python 3.10+,至少 4C8G 配置 |
| 显存需求 | 联邦搜索本身不需要 GPU;如果联动本地 embedding 模型则按模型需求评估 |
| 启动方式 | 服务化部署或作为 Python 库嵌入 Agent 进程 |
| 适用场景 | 企业知识库问答、多系统数据聚合、复杂 Agent 工具编排、统一检索网关 |
说明一下,Federated Search 并不是新概念,在搜索引擎和企业信息检索领域已经存在很多年。这里的关键变化是:它被重新放到“AI Agent 的工具链”里,需要处理的是自然语言查询、异构 API、非结构化数据和 RAG 检索链路之间的适配问题。因此后面所有讨论都围绕“Agent 场景下的联邦搜索”展开。
2. 适用场景与使用边界
2.1 适合解决的典型问题
第一类场景是 RAG 多库检索。传统 RAG 只挂一个向量库,遇到“哪些客户在最近一个月提过退款诉求且订单金额大于 5000”这种组合查询时完全不够用。联邦搜索可以把向量库、订单数据库、工单 API 同时纳入查询范围,让 Agent 拿到跨源聚合后的结果再作答。
第二类场景是多工具编排。Agent 经常需要调用多个外部工具,比如机票查询、酒店查询、天气查询。如果没有统一检索层,Agent 要自己在代码里维护多路调用、超时处理和结果合并逻辑。联邦搜索把这个过程收敛成一个 Tool,Agent 只负责生成查询参数和消费最终结果。
第三类场景是企业统一搜索网关。很多公司内部同时存在 Confluence、Jira、飞书文档、内部 Wiki、代码仓库等系统。联邦搜索可以屏蔽掉每个系统的访问差异,提供一个统一检索接口给上层 Agent,同时在入口统一做权限校验、审计日志和限流。
2.2 不适合什么场景
如果只有一个数据源,或者数据已经全部同步到一个统一索引里,联邦搜索就是多余的。它最大的价值在于处理“多源异构”和“并发聚合”,单一数据源场景只会增加架构复杂度。
另外,如果数据源之间没有任何关联字段,且查询结果不需要合并,那也不用上联邦搜索。这时候直接让 Agent 依次调用多个工具可能更简单。
2.3 合规与安全边界
联邦搜索会把请求分发到多个系统,再把结果聚回来,这里有几条必须注意的红线:
- 必须有统一的访问控制和数据权限校验,不能因为联邦层做了聚合就绕过单系统的权限。
- 涉及用户隐私、人脸、身份证、手机号等敏感字段时,返回给 Agent 之前要做字段级脱敏。
- 如果数据源包含版权材料或未授权的内部资料,要确认检索和传播边界。
- 记录完整的审计日志,尤其是 Agent 自动化调用时,需要能回溯“谁在什么时间查了什么数据”。
- 在公开网络环境部署时,务必限制联邦搜索服务的访问范围,避免接口被恶意批量调用。
3. 联邦搜索的架构拆解
3.1 整体架构
一个面向 Agent 的联邦搜索系统通常包含五个核心模块。
用户请求 / Agent 调用 ↓ 查询规划器(Query Planner) ↓ 联邦分发器(Federated Dispatcher) ↓ 数据源连接器(Source Connector)→ 数据源 A 数据源连接器(Source Connector)→ 数据源 B 数据源连接器(Source Connector)→ 数据源 C ↓ 结果处理器(Result Processor) ↓ 响应合成器(Response Synthesizer) ↓ Agent / 用户查询规划器负责把自然语言查询或者结构化参数转换成每个数据源能接受的查询条件。联邦分发器负责并发调度,决定哪些数据源需要查询、哪些可以并行、哪些必须串行。数据源连接器是适配层,屏蔽不同系统的接口差异。结果处理器负责去重、合并、评分和截断。响应合成器把多源结果整理成统一结构返回给上层。
3.2 查询规划的关键点
联邦搜索里最核心的是查询规划。一个常见做法是把查询拆分为多个子查询,每个子查询对应一个数据源。
比如 Agent 收到请求“总结一下最近一周技术部提交的高优先级工单情况”,查询规划器需要拆成:
- 工单系统查询:最近一周、技术部、高优先级
- 人员系统查询:技术部成员列表或部门 ID
- 文档系统查询:这些工单关联的团队文档或复盘记录
这里最常用的规划策略是基于数据源 schema 的任务分解。每个数据源在配置里声明自己能处理哪些查询字段,规划器根据查询语义把参数映射到对应字段上。
3.3 数据源连接器的适配层级
连接器通常做三件事:协议转换、参数映射、结果标准化。协议转换解决的是“这个系统是 REST 接口,那个系统是 SQL,另一个是 SDK 调用”的问题;参数映射解决的是“同一个概念在不同系统里字段名不一样”的问题;结果标准化则是把各种数据结构统一成 Schema,方便后续聚合。
从工程实践看,连接器不要手写太多。优先使用现成 SDK,或者用声明式配置描述数据源的请求格式、认证方式和返回字段映射,尽量减少硬编码。
4. 环境准备与前置条件
4.1 基础运行环境
联邦搜索本身不需要 GPU。它更依赖网络 IO 和并发调度能力。建议按以下清单准备环境:
| 检查项 | 建议要求 |
|---|---|
| 操作系统 | Linux(Ubuntu 20.04+)或 macOS,Windows 可用于本地调试 |
| Python | 3.10 或更高版本 |
| 依赖管理 | pip + virtualenv 或 uv,服务化部署时建议用 venv 隔离 |
| CPU / 内存 | 4 核 + 8G 内存起步,多数据源并发场景建议 8 核 + 16G |
| 网络 | 能访问待查询的所有数据源服务 |
| 磁盘 | 根据查询日志和缓存策略评估,一般预留 10G 足够 |
| 服务端口 | 预留一个 HTTP 服务端口,例如 8080 或 8900 |
需要注意,如果联邦搜索在 AI Agent 的生成链路里,并且要联动 embedding 模型、rerank 模型,那显存占用就是另一回事了。embedding 和 rerank 模型通常跑在 GPU 上,显存需求取决于模型体积。6G 显存可以运行中小规模 embedding 模型,更大模型需要按实际情况评估。联邦搜索网关本身不做向量计算,不要把它和模型推理混在一起做性能评估。
4.2 通用依赖安装示例
下面给出一套通用的 Python 项目依赖安装示例。这里没有针对某个具体开源项目,而是让你快速搭起一个联邦搜索原型环境:
mkdir federated-search-agent cd federated-search-agent python3.10 -m venv venv source venv/bin/activate # 核心依赖 pip install fastapi uvicorn httpx pydantic python-dotenv # 如果涉及数据库查询 pip install sqlalchemy psycopg2-binary # 如果涉及向量数据库 pip install pymilvus chromadb # 如果涉及 Elasticsearch pip install elasticsearch安装完成后可以验证一下核心依赖是否可导入:
python -c "import fastapi, httpx, pydantic; print('deps ok')"4.3 确认数据源连通性
在开发联邦搜索之前,先确认每个数据源 API 能通。这一步很多人会忽略,结果系统写完才发现某个数据源连接串配置错误。
建议准备一个 init_check 脚本,把数据源连通性检查放到启动项里,任何必要数据源不可达时直接拒绝启动。
5. 一个最小可运行的联邦搜索调度设计
下面给出一个简化的 Python 联邦搜索调度示例。它不依赖任何重型框架,只用来演示核心流程:把查询分发到多个数据源,等待结果,再对结果做合并。
5.1 数据源配置
我们用 YAML 或 JSON 来描述数据源。每个数据源包含名称、类型、endpoint 和超时时间。
sources: wiki: type: http endpoint: "https://your-wiki-endpoint/api/search" timeout_seconds: 5 ticket: type: http endpoint: "https://your-ticket-system/api/tickets/search" timeout_seconds: 8 vector: type: http endpoint: "http://127.0.0.1:8001/query" timeout_seconds: 3严格来说,这里的 http 类型连接器就是一个通用请求封装,生产环境里每个数据源都会有独立的参数映射逻辑。
5.2 调度器代码
import asyncio import httpx from typing import Any, Dict, List async def query_source( client: httpx.AsyncClient, source_name: str, source_config: Dict[str, Any], query_payload: Dict[str, Any], ) -> Dict[str, Any]: """向单个数据源发起查询,返回统一格式的结果。""" try: endpoint = source_config["endpoint"] timeout = source_config.get("timeout_seconds", 5) response = await client.post( endpoint, json=query_payload, timeout=timeout, ) response.raise_for_status() data = response.json() return { "source": source_name, "success": True, "items": data.get("items", data.get("results", [])), } except Exception as exc: return { "source": source_name, "success": False, "error": str(exc), "items": [], } async def federated_search( source_configs: Dict[str, Dict[str, Any]], query_payloads: Dict[str, Dict[str, Any]], ) -> List[Dict[str, Any]]: """并发查询多个数据源并聚合结果。""" async with httpx.AsyncClient() as client: tasks = [ query_source(client, name, config, query_payloads[name]) for name, config in source_configs.items() if name in query_payloads ] results = await asyncio.gather(*tasks, return_exceptions=False) return results if __name__ == "__main__": # 演示配置 configs = { "wiki": { "endpoint": "http://127.0.0.1:9001/search", "timeout_seconds": 5, }, "ticket": { "endpoint": "http://127.0.0.1:9002/search", "timeout_seconds": 8, }, } payloads = { "wiki": {"query": "最近一周技术部文档", "top_k": 10}, "ticket": {"query": "高优先级工单", "status": "open"}, } merged = asyncio.run(federated_search(configs, payloads)) for item in merged: print(item["source"], item["success"], len(item["items"]))这是一个很基础的并发查询骨架。实际项目中还需要加入查询规划、鉴权、重试、超时降级、结果去重重排等能力。
5.3 启动一个 HTTP 服务
为了让 Agent 能方便地调用,通常会把联邦搜索封装成 HTTP 服务。下面用 FastAPI 做一个最小示例:
from fastapi import FastAPI from pydantic import BaseModel from typing import Dict, Any app = FastAPI(title="Federated Search For AI Agents") class SearchRequest(BaseModel): query: str sources: list[str] = [] top_k: int = 10 class SearchResponse(BaseModel): query: str results: list[Dict[str, Any]] source_stats: Dict[str, Any] @app.post("/search", response_model=SearchResponse) async def search(req: SearchRequest): # 这里应该接入真实的数据源配置和调度器 # 示例只做流程演示 results = [ {"source": "demo", "score": 0.9, "content": "演示结果"} ] return SearchResponse( query=req.query, results=results, source_stats={"demo": {"count": 1, "elapsed_ms": 42}}, ) if __name__ == "__main__": import uvicorn uvicorn.run(app, host="127.0.0.1", port=8080)启动方式:
python serve.py启动后,Agent 或者人工测试可以直接请求:
curl -X POST http://127.0.0.1:8080/search \ -H "Content-Type: application/json" \ -d '{"query": "最近一周工单情况", "sources": ["ticket", "wiki"], "top_k": 5}'这里的查询参数、返回结构都是示例,实际需要根据项目接口调整。
6. 功能测试与效果验证
6.1 单数据源连通性测试
先测单个数据源,确保连接器和请求映射正确。
测试步骤如下:
- 准备一个已知答案的查询。比如数据源里有一条确定存在的工单,查询它的标题。
- 调用联邦搜索服务的 /search 接口,只指定一个 source。
- 检查返回结果中是否包含预期记录。
- 人为把 endpoint 改成错误地址,确认返回 error 结构能正确透出。
判断标准:指定单源查询时,返回结果与该数据源单独提供的结果一致。
6.2 多数据源合并测试
多源测试要关注三个点:结果是否都返回、重复项是否被处理、字段冲突如何解决。
比如同时查向量库和 Wiki 文档库,两者可能返回同一篇文档的不同版本。好一点的联邦搜索应该有字段级合并策略:以 ID 作为去重键,保留更新时间最新的版本,或者保留 score 更高的版本。
更稳妥的测试方式是构造一组数据,让两个数据源都返回同一条记录,然后在输出中检查该记录只出现一次,并且来源字段展示了多源命中情况。
6.3 查询延迟与超时降级测试
联邦搜索的最大风险是一个慢数据源拖垮整个链路。所以要把超时控制当成功能来测。
测试方法:
- 将一个数据源接口调成 30 秒后返回。
- 将连接器超时设置为 5 秒。
- 发起查询,观察整个请求是否在 5 秒左右返回。
- 检查返回结果里该数据源是否标记为 failed,并携带超时错误信息。
推荐在系统设计上做到:单个数据源失败不影响整体结果返回。降级后的结果可以缺失部分来源,但接口仍应该返回 200,并在响应里标注哪些数据源查询失败。
6.4 Agent 集成测试
最后把联邦搜索接入 Agent 的工具调用链路。这里测试的是“Agent 是否会正确生成查询参数,并消费返回结果”。
操作步骤:
- 在 Agent 的 Tool 配置里注册联邦搜索服务。
- 给 Agent 一个复杂的跨源问题,例如“对比技术部最近一个月和上个月的工单解决率”。
- 检查 Agent 是否调用了联邦搜索接口。
- 检查返回结果是否被大模型正确引用。
- 检查最终回答是否包含了多个数据源的信息。
判断标准:Agent 生成的查询参数符合预期,联邦搜索结果完整传回,大模型基于多源结果给出了可验证的回答。
7. 接口 API 与批量任务设计
7.1 API 路径设计
面向 Agent 的联邦搜索接口一般会提供以下端点:
| 接口 | 方法 | 用途 |
|---|---|---|
| /health | GET | 健康检查 |
| /search | POST | 单次联邦搜索 |
| /batch_search | POST | 批量搜索任务提交 |
| /batch/{task_id} | GET | 查询批量任务状态和结果 |
| /sources | GET | 获取可用数据源列表 |
7.2 批量任务设计
批量场景里常见的问题是:短时间内大量 Agent 请求同时进来,不能直接同步循环调用,否则会把数据源打挂。
推荐用异步任务队列。思路是:
- 提交批量任务:接收一组查询请求,生成一个 task_id。
- 后台消费者从队列里取出请求,逐个调用联邦搜索。
- 每个查询独立记录状态,失败自动重试。
- 调用方通过 task_id 查询进度和结果。
示例请求格式:
{ "task_name": "batch_ticket_analysis", "queries": [ { "query": "高优先级工单", "sources": ["ticket", "wiki"], "top_k": 5 }, { "query": "低优先级工单", "sources": ["ticket"], "top_k": 3 } ] }Python 批量调用示例:
import requests import time url = "http://127.0.0.1:8080/batch_search" payload = { "task_name": "demo_batch", "queries": [ {"query": "工单A", "sources": ["ticket"], "top_k": 3}, {"query": "工单B", "sources": ["ticket", "wiki"], "top_k": 5}, ], } response = requests.post(url, json=payload, timeout=10) task_id = response.json().get("task_id") print("task_id:", task_id) # 轮询任务结果 status_url = f"http://127.0.0.1:8080/batch/{task_id}" for _ in range(30): resp = requests.get(status_url, timeout=10) data = resp.json() if data["status"] in ("completed", "failed"): print(data) break time.sleep(2)需要注意,这里的接口路径和返回结构是通用示例,实际需要按项目具体实现调整。
7.3 重试与失败隔离
批量任务里最容易踩的坑是数据源偶尔超时。如果重试逻辑写得太激进,反而会放大流量压力。
比较稳妥的重试策略是:
- 连接超时短一点,读超时长一点。
- 对于超时失败的请求,最多重试 1 到 2 次。
- 使用指数退避,第一次等 1 秒,第二次等 2 秒。
- 数据源连续失败超过阈值时,直接熔断,本轮任务标记失败,不再继续重试。
另外,要给每个查询设置独立超时上下文。避免一个数据源卡住导致整个任务队列堆积。
8. 资源占用与性能观察
8.1 观察哪些指标
联邦搜索的资源占用核心不在 CPU 和内存,而在网络连接、并发线程数、队列积压和下游数据源负载。
建议采集以下指标:
| 指标 | 说明 |
|---|---|
| QPS | 每秒钟处理的联邦搜索请求数 |
| 平均耗时 / P95 耗时 | 请求从进入到返回的时间分布 |
| 各数据源耗时占比 | 哪个数据源最慢 |
| 超时数量 | 每个时间窗口内超时的查询数量 |
| 队列长度 | 批量任务排队数量 |
| 下游连接池占用 | 到每个数据源的连接数是否打满 |
8.2 通过日志定位慢数据源
在代码里给每个数据源查询打上耗时标签:
source=ticket elapsed_ms=3200 status=success source=wiki elapsed_ms=150 status=success source=vector elapsed_ms=8000 status=timeout这样一眼就能看出哪个环节是瓶颈。如果某个数据源经常慢,可以考虑做两级缓存,或者把这个数据源设置成异步刷新,不阻塞主链路。
8.3 降低资源占用的策略
- 连接复用:使用 httpx.AsyncClient 时复用 client 实例,不要每次查询都新建连接。
- 结果截断:各数据源返回的 top_k 数量要控制住,不要让超大数据集进入聚合层。
- 缓存命中:对高频查询做短 TTL 缓存,比如 30 到 60 秒。
- 并发限流:给每个数据源配置最大并发数,防止瞬时流量打爆下游。
- 降级开关:某个数据源故障时,可以动态把该源从联邦搜索里摘掉。
8.4 本地性能测试的通用流程
本地验证时,可以先用 mock 数据源测试调度器本身的性能。mock 数据源延迟设置为 20ms,然后用并发请求打压联邦搜索服务,观察吞吐量和延迟。这一步能帮你判断调度器本身有没有明显问题,而不用先背真实数据源的锅。
9. 常见问题与排查方法
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 某个数据源结果一直为空 | 数据源连接串错误、参数映射不对、权限不足 | 单源模式单独调试,查看返回原始报文 | 先用 curl 测原始接口,再对比连接器映射 |
| 查询偶尔超时 | 下游数据源慢、超时时间设置过短 | 看日志里耗时分布和超时记录 | 调整超时参数,增加重试和降级策略 |
| 结果重复严重 | 缺少去重逻辑,或者去重键设置不合理 | 检查两个数据源返回记录 ID 是否一致 | 在结果处理器里按业务主键去重 |
| 字段含义冲突 | 不同数据源对同一字段定义不同 | 查看 Schema 映射文档 | 在连接器层做字段名标准化,或在配置中声明优先级 |
| Agent 拿到结果后回答错误 | 返回给 Agent 的结果结构不清晰,缺失来源标注 | 检查 Agent 工具返回的 prompt 上下文 | 在结果里补充 source 字段和 score 字段,让模型知道信息来源 |
| 批量任务堆积 | 下游数据源并发能力不足,队列消费不过来 | 看队列长度和下游连接池监控 | 增加消费者线程数,或降低峰值并发 |
| 启动后服务端口被占用 | 端口冲突 | 查看启动日志 | 换一个端口,例如 8081 |
| 依赖安装失败 | Python 版本不兼容或缺少编译工具 | 查看 pip 安装日志 | 用更高版本 Python,或者使用 conda 环境 |
10. 最佳实践与使用建议
10.1 先小规模跑通,再上复杂场景
第一次做联邦搜索,建议只接两个数据源:一个文档类、一个关系型数据库。先跑通一个跨源查询,确认查询规划、结果合并逻辑能正常工作,再逐步增加数据源。
很多团队一上来就接七八个系统,结果问题全出在连接器和权限上,核心的查询规划反而没验证到。
10.2 配置和代码分离,数据源声明式管理
数据源配置不要硬编码到代码里。用 YAML、JSON 或者配置中心来管理。每个数据源独立配置认证方式、超时时间、并发限制和字段映射。这样后续调整某个系统的连接参数时,不需要重新发版。
10.3 保留审计日志
Agent 自动调用会放大查询量。建议记录每次查询的输入、命中的数据源、返回的结果条数和耗时。一旦出现数据泄露或者异常调用,审计日志是最重要的回溯依据。
10.4 给 Agent 一个“结构化结果”
联邦搜索的最终输出,要和普通搜索引擎区分开。Agent 消费的结果最好是结构化数据,带着字段名,而不是一长串拼接后的文本。例如:
{ "source": "ticket", "id": "TICKET-1234", "title": "登录模块偶发超时", "priority": "high", "status": "open" }清晰的字段结构能让 Agent 在生成回答时更稳定,也能减少模型误读。
10.5 权限收敛在联邦层,也在数据源层
联邦搜索是统一的入口,但每个数据源的权限校验不能省。正确的做法是联邦层做身份认证和范围限制,数据源层做最终的数据权限校验,两层同时生效。
10.6 内容合规与授权提醒
如果联邦搜索接入的是外部数据源或第三方接口,需要确认数据使用合法性。对涉及人脸、声音、隐私信息、版权文本的内容,必须在检索和聚合阶段做权限校验、错误处理和字段脱敏。不要在未获授权的情况下,将来源系统的内部数据作为 Agent 检索结果转发或商用。
11. 总结与下一步
Federated Search for AI Agents 最有价值的地方,是把 Agent 从“单库检索”推进到“多源统一检索”。它不解决模型能力问题,而是解决数据通路问题。对做 RAG、企业知识库、Agent 工具编排的开发者来说,这套架构能明显提升检索覆盖率和回答可信度。
建议最先验证的能力是:两个不同数据源能否通过一个接口并行查询,并正确合并结果。最容易踩的坑是超时控制和结果去重——一个慢数据源拖垮全链路、两个重复记录干扰模型回答,这两类问题在真实场景基本必然遇到。
后续可以继续扩展的方向有三个:一是把查询规划做成基于 LLM 的自动路由,让模型决定哪些数据源参与本次查询;二是增加 rerank 层,统一对多源结果做质量打分;三是把联邦搜索抽象成 Agent 标准工具,接入 LangChain、Semantic Kernel 等 Agent 框架。先从最小架构跑通,再逐步迭代,这条路是稳的。