行业资讯
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: OptionString, } /// 工具参数定义 #[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: VecToolParam, } /// Agent 执行的单个步骤 #[derive(Debug)] pub enum AgentStep { /// 调用某个工具并传入参数 ToolCall { tool_name: String, arguments: HashMapString, 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: HashMapString, serde_json::Value, ) - ResultToolResult, 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: HashMapString, serde_json::Value, ) - ResultToolResult, 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: Vecsqlx::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: Vecsqlx::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: HashMapString, serde_json::Value, ) - ResultToolResult, 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: VecBoxdyn AgentTool, /// 本地 AI 模型客户端 llm: LlmClient, } impl AgentOrchestrator { /// 执行一次 Agent 任务 pub async fn run(self, user_query: str) - ResultString, String { let mut collected_data: VecToolResult Vec::new(); let mut max_rounds 5; // 最多 5 轮工具调用防止无限循环 // 构建工具列表描述发送给 AI 做规划 let tool_descriptions: VecToolDef 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], ) - ResultAgentStep, 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: HashMapString, 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], ) - ResultString, 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 时代软件工程的核心范式转变。有问题欢迎评论区讨论
郑州网站建设
网页设计
企业官网