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 对零结果做二次确认。

五、总结

  1. 工具注册机制:通过 AgentTool trait 统一不同数据源的接口,增删数据源只需实现 trait,不影响调度逻辑。
  2. AI 决策编排:Agent 把用户问题、可用工具、已收集数据发给 LLM,让 AI 动态决定"下一步该查什么"。
  3. 多轮迭代收集:支持最多 5 轮工具调用,AI 可以在看到前期结果后调整后续查询方向。
  4. 自动报告生成:所有数据收集完毕后,AI 自动汇总成 Markdown 格式的分析报告。
  5. 安全性:SQL 工具仅允许 SELECT 查询,HTTP 工具限制域名白名单(可扩展),防止越权操作。

这套方案的生产化方向:

  • 添加工具调用缓存,相同查询短时间内直接复用
  • 添加参数化报告模板,按场景生成不同格式(PPT、邮件、飞书文档)
  • 添加人工确认环节,关键步骤(如发送通知)需要用户二次确认

Agent 的精髓在于"让 AI 自己决定下一步做什么",而不是死板的 if-else 逻辑。这种灵活性正是 AI 时代软件工程的核心范式转变。

有问题欢迎评论区讨论!

Logo

中国智能体开发者社区,聚焦智能体与大模型开发,提供前沿资讯、实用工具链、开源项目及行业案例。通过技术沙龙、开发者大赛等活动,促进经验交流与协作,助力开发者快速构建创新智能应用。

更多推荐