AI Agent 数据聚合:多个数据源查询结果,让模型自动汇总成报告
AI Agent 数据聚合:多个数据源查询结果,让模型自动汇总成报告
一、场景:信息整合的噩梦
大家好,我是一铭。上个月领导让我做一份竞品分析报告,我需要从五个数据源收集信息:公司的内部数据库、第三方 API、网页爬虫数据、用户反馈系统、搜索引擎结果。每个数据源返回的数据格式都不一样,SQL 结果、JSON API、HTML 页面、文本反馈……
以前我的做法是:分别查,分别整理到 Excel,然后人工写报告。一整天就过去了。
有没有自动化方案?让 AI Agent 自动查询多个数据源,然后把结果汇总成一份结构化的报告?
答案是有的,而且用 Rust + AI 可以做得非常丝滑。
二、Agent 架构设计
2.1 核心组件
一个 AI Agent 数据聚合系统包含三个核心组件:
- Tool 注册表:定义每个数据源的工具(名称、描述、参数、执行函数)
- Agent 调度器:接收用户问题,调用 AI 规划查询步骤,按序执行工具
- 报告生成器:收集所有工具返回的数据,调用 AI 生成汇总报告
2.2 Rust 类型定义
use serde::{Deserialize, Serialize}; use async_trait::async_trait; use std::collections::HashMap; /// 工具执行结果 #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ToolResult { /// 工具名称 pub tool_name: String, /// 执行是否成功 pub success: bool, /// 返回的数据(任意 JSON 结构) pub data: serde_json::Value, /// 执行耗时(毫秒) pub elapsed_ms: u64, /// 错误信息(如果有) pub error: Option<String>, } /// 工具参数定义 #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ToolParam { /// 参数名 pub name: String, /// 参数描述(给 AI 看的,用于理解参数的用途) pub description: String, /// 参数类型 pub param_type: String, /// 是否必填 pub required: bool, } /// 工具定义:告诉 AI 这个工具能做什么 #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ToolDef { /// 工具名称,AI 通过名称来调用 pub name: String, /// 工具描述:详细说明功能和适用场景 pub description: String, /// 参数列表 pub params: Vec<ToolParam>, } /// Agent 执行的单个步骤 #[derive(Debug)] pub enum AgentStep { /// 调用某个工具,并传入参数 ToolCall { tool_name: String, arguments: HashMap<String, serde_json::Value>, }, /// 将前面收集到的所有数据汇总为报告 Summarize, /// 任务完成 Done { report: String }, }2.3 工具注册与执行
use async_trait::async_trait; /// 每个工具都必须实现这个 trait #[async_trait] pub trait AgentTool: Send + Sync { /// 返回工具的元数据定义 fn definition(&self) -> ToolDef; /// 执行工具,返回结果 async fn execute( &self, args: HashMap<String, serde_json::Value>, ) -> Result<ToolResult, String>; } // ---- 实现具体的工具 ---- /// SQL 数据库查询工具 pub struct SqlQueryTool { db_pool: sqlx::PgPool, } #[async_trait] impl AgentTool for SqlQueryTool { fn definition(&self) -> ToolDef { ToolDef { name: "sql_query".into(), description: "查询内部 PostgreSQL 数据库,返回结构化数据。用于获取用户信息、订单记录、业务指标等。".into(), params: vec![ToolParam { name: "query".into(), description: "要执行的 SQL 查询语句(仅支持 SELECT)".into(), param_type: "string".into(), required: true, }], } } async fn execute( &self, args: HashMap<String, serde_json::Value>, ) -> Result<ToolResult, String> { let start = std::time::Instant::now(); let query = args.get("query") .and_then(|v| v.as_str()) .ok_or("缺少 query 参数")?; // 安全:只允许 SELECT 查询,防止 SQL 注入 let trimmed = query.trim().to_uppercase(); if !trimmed.starts_with("SELECT") { return Err("仅允许 SELECT 查询".into()); } // 执行查询 let rows: Vec<sqlx::postgres::PgRow> = sqlx::query(query) .fetch_all(&self.db_pool) .await .map_err(|e| format!("SQL 执行失败: {}", e))?; // 转换为 JSON let data = rows_to_json(rows); Ok(ToolResult { tool_name: "sql_query".into(), success: true, data, elapsed_ms: start.elapsed().as_millis() as u64, error: None, }) } } // 辅助函数:将查询行转为 JSON fn rows_to_json(rows: Vec<sqlx::postgres::PgRow>) -> serde_json::Value { serde_json::Value::Array( rows.iter().map(|row| { // 简化:实际应该遍历列并转换类型 serde_json::json!({"row": format!("{:?}", row)}) }).collect() ) } /// HTTP API 调用工具 pub struct HttpApiTool; #[async_trait] impl AgentTool for HttpApiTool { fn definition(&self) -> ToolDef { ToolDef { name: "http_api".into(), description: "调用外部 HTTP API 获取数据。用于获取第三方服务的数据,如天气、股票、新闻等。".into(), params: vec![ ToolParam { name: "url".into(), description: "API 地址".into(), param_type: "string".into(), required: true, }, ToolParam { name: "method".into(), description: "HTTP 方法(GET/POST)".into(), param_type: "string".into(), required: false, }, ], } } async fn execute( &self, args: HashMap<String, serde_json::Value>, ) -> Result<ToolResult, String> { let start = std::time::Instant::now(); let url = args.get("url") .and_then(|v| v.as_str()) .ok_or("缺少 url 参数")?; let method = args.get("method") .and_then(|v| v.as_str()) .unwrap_or("GET"); let client = reqwest::Client::new(); let response = match method { "GET" => client.get(url).send().await, "POST" => client.post(url).send().await, _ => return Err(format!("不支持的 HTTP 方法: {}", method)), } .map_err(|e| format!("HTTP 请求失败: {}", e))?; let body = response .text() .await .map_err(|e| format!("读取响应失败: {}", e))?; // 尝试解析为 JSON,如果失败就保留原文 let data = serde_json::from_str::<serde_json::Value>(&body) .unwrap_or(serde_json::Value::String(body)); Ok(ToolResult { tool_name: "http_api".into(), success: true, data, elapsed_ms: start.elapsed().as_millis() as u64, error: None, }) } }三、Agent 调度器:让 AI 编排工具调用
/// AI Agent 调度器 pub struct AgentOrchestrator { /// 注册的工具列表 tools: Vec<Box<dyn AgentTool>>, /// 本地 AI 模型客户端 llm: LlmClient, } impl AgentOrchestrator { /// 执行一次 Agent 任务 pub async fn run(&self, user_query: &str) -> Result<String, String> { let mut collected_data: Vec<ToolResult> = Vec::new(); let mut max_rounds = 5; // 最多 5 轮工具调用,防止无限循环 // 构建工具列表描述,发送给 AI 做规划 let tool_descriptions: Vec<ToolDef> = self.tools .iter() .map(|t| t.definition()) .collect(); while max_rounds > 0 { max_rounds -= 1; // 第一步:让 AI 决定下一步做什么 let step = self.plan_next_step( user_query, &tool_descriptions, &collected_data, ).await?; match step { AgentStep::ToolCall { tool_name, arguments } => { // 找到对应的工具并执行 let tool = self.tools.iter() .find(|t| t.definition().name == tool_name) .ok_or(format!("未找到工具: {}", tool_name))?; let result = tool.execute(arguments).await?; collected_data.push(result); } AgentStep::Summarize => { // 让 AI 生成最终报告 let report = self.generate_report( user_query, &collected_data, ).await?; return Ok(report); } AgentStep::Done { report } => { return Ok(report); } } } Err("Agent 达到最大轮次限制".into()) } /// 让 AI 规划下一步 async fn plan_next_step( &self, query: &str, tools: &[ToolDef], history: &[ToolResult], ) -> Result<AgentStep, String> { // 构建 prompt:告诉 AI 它有哪些工具可用,目前收集了哪些数据 let tools_json = serde_json::to_string_pretty(tools).unwrap_or_default(); let history_json = serde_json::to_string_pretty(history).unwrap_or_default(); let prompt = format!( r#"你是一个数据聚合 Agent。用户的问题是: "{query}" 你可以使用以下工具: {tools_json} 已收集到的数据: {history_json} 请返回 JSON 格式的下一步计划。如果还需要收集数据,返回: {{"action": "tool_call", "tool": "工具名", "args": {{...}}}} 如果数据足够生成报告,返回: {{"action": "summarize"}} 只返回 JSON,不要其他内容。"# ); // 调用 LLM 获取计划(此处简化) let response = self.llm.chat(&prompt).await?; let plan: serde_json::Value = serde_json::from_str(&response) .map_err(|e| format!("AI 返回无效 JSON: {}", e))?; match plan["action"].as_str() { Some("tool_call") => { let tool_name = plan["tool"].as_str().unwrap_or("").to_string(); let args: HashMap<String, serde_json::Value> = plan["args"] .as_object() .map(|obj| obj.iter() .map(|(k, v)| (k.clone(), v.clone())) .collect() ) .unwrap_or_default(); Ok(AgentStep::ToolCall { tool_name, arguments: args }) } Some("summarize") => Ok(AgentStep::Summarize), _ => Err("AI 返回了无法识别的动作".into()), } } /// 生成最终报告 async fn generate_report( &self, query: &str, results: &[ToolResult], ) -> Result<String, String> { let data_json = serde_json::to_string_pretty(results).unwrap_or_default(); let prompt = format!( r#"用户问题:"{query}" 从多个数据源收集到的数据: {data_json} 请基于以上数据生成一份结构化的分析报告,使用 Markdown 格式。 要求: 1. 包含"数据概览"一节,总结数据来源和关键指标 2. 包含"详细分析"一节,深入解读数据 3. 包含"结论与建议"一节 4. 语言专业、客观"# ); self.llm.chat(&prompt).await } }四、完整执行流程
实战踩坑:工具超时与部分数据聚合
上了这套 Agent 后不久,遇到了两个现实问题:
问题一:工具超时导致整体任务失败。有次 HTTP API 工具请求第三方服务卡了 30 秒,整个 Agent 就挂住了。因为调度器的max_rounds循环是串行的——一个工具超时,后面所有工具都等不到。解法:给每个工具执行加超时:
// 在 AgentStep::ToolCall 分支中加超时保护 let result = tokio::time::timeout( Duration::from_secs(10), tool.execute(arguments), ).await .map_err(|_| "工具执行超时".to_string())?;问题二:AI 生成不存在的工具参数。Agent 框架告诉 AI 工具需要query参数,但 AI 有时会编一个sql参数名(因为它对 SQL 更熟悉)。导致执行时报"缺少 query 参数"。解决办法是增加一层参数校验,如果 AI 传了不认识的参数名,自动做 fuzzy matching 尝试纠错,或者在 prompt 里强调参数名必须精确匹配。
还有一个认知收获:不要让 Agent 在发现部分工具失败后仍然尝试汇总。汇总出的报告会基于不完整的数据,结论可能有误导性。更好的做法是返回"部分数据源不可用,报告仅基于 X/Y 个数据源",并标红缺失的部分。
实际项目里还发现一个"Agent 过度自信"的问题:SQL 工具查出了空结果,Agent 直接报告"数据库中没有相关数据"——但其实是因为查询的表名写错了。后来在所有工具的返回结果里加了result_count字段,并在汇总阶段让 AI 对零结果做二次确认。
五、总结
- 工具注册机制:通过
AgentTooltrait 统一不同数据源的接口,增删数据源只需实现 trait,不影响调度逻辑。 - AI 决策编排:Agent 把用户问题、可用工具、已收集数据发给 LLM,让 AI 动态决定"下一步该查什么"。
- 多轮迭代收集:支持最多 5 轮工具调用,AI 可以在看到前期结果后调整后续查询方向。
- 自动报告生成:所有数据收集完毕后,AI 自动汇总成 Markdown 格式的分析报告。
- 安全性:SQL 工具仅允许 SELECT 查询,HTTP 工具限制域名白名单(可扩展),防止越权操作。
这套方案的生产化方向:
- 添加工具调用缓存,相同查询短时间内直接复用
- 添加参数化报告模板,按场景生成不同格式(PPT、邮件、飞书文档)
- 添加人工确认环节,关键步骤(如发送通知)需要用户二次确认
Agent 的精髓在于"让 AI 自己决定下一步做什么",而不是死板的 if-else 逻辑。这种灵活性正是 AI 时代软件工程的核心范式转变。
有问题欢迎评论区讨论!
