Agent 开发实战(二):对接多源后台数据的分析 Agent

Moryu 3 阅读 agent

Agent 开发实战(二):对接多源后台数据的分析 Agent

真实网站的"后台数据"从来不止一张数据库表。它可能是 MySQL 里的业务表、Excel/CSV 报表、产品/政策文档(PDF/Word/Markdown),也可能是某个垂直系统特有的文件——比如生信分析系统导出的 VCF 变异文件、FASTA 序列、表达矩阵等。本文手把手实现一个统一接入多源数据、用 LLM 做跨源分析的 Agent,并重点讲清楚每类数据源的安全接入方式。

这是「Agent 开发实战」系列的第 2 篇,承接(一)跑通的最小 Agent。(一)已经把"模型决策 / 工具执行"的核心循环讲清楚了,本篇直接复用这个模式,重点放在把多种数据源变成受控工具、让 Agent 能回答跨源问题。后续(三)~(九)会依次深入工具编排、记忆、前端、RAG、多 Agent、可观测性与工程化部署。

系列地图:(一)最小 Agent | (二)多源数据接入 ← 本篇 | (三)工具系统深入 | (四)记忆与上下文 | (五)流式与前端集成 | (六)RAG 进阶 | (七)多 Agent 协作 | (八)可观测性与评估 | (九)工程化部署与端到端案例。逻辑是「地基 → 接数据 → 工具工程化 → 记忆 → 前端 → 检索 → 多 Agent → 观测 → 部署」,每篇都基于前一篇叠加一层能力,但共享同一套"模型决策、工具执行、每源护栏"的底座。

1. 我们要解决什么

一个典型网站后台,数据散落在不同形态里:

  • 数据库usersordersposts 等结构化业务表
  • 表格数据:运营导出的 sales_2026Q3.xlsx、活动报名 signup.csv
  • 文档:《退款政策》《产品白皮书》这类 PDF / Word / Markdown
  • 系统特有文件:生信分析系统的 sample1.vcf.gz(变异)、genome.fasta(序列)、expression.tsv(表达矩阵)

业务问题往往是跨源的:"上季度 A 类商品退货偏高,结合数据库订单、《退款政策》文档,分析退货集中在哪些渠道;再调出生信样本里是否有相关变异标记。"——这类问题靠人肉拼数据非常痛苦,正是 Agent 的用武之地。

Agent 的价值:把多源数据变成统一的自然语言入口,让模型自己决定去哪个源取数、怎么交叉分析。

2. 架构设计

┌──────────┐   自然语言    ┌───────────────────────────────┐
│  用户     │ ──────────▶ │      Analysis Agent (Python)    │
│ (运营/你) │             │  ├─ LLM (function calling)      │
└──────────┘ ◀────────── │  ├─ Tool 层(统一接口)          │
       ▲      分析结论     │  └─ 安全护栏(每源独立约束)      │
       │                  └───────┬─────────┬────────┬──────┘
       │                          │         │        │
       │                  ┌───────┴──┐ ┌────┴───┐ ┌──┴─────────┐
       │                  │ 数据库    │ │ 表格/文档│ │ 系统特有文件 │
       │                  │ MySQL RO │ │ Excel/  │ │ VCF/FASTA/  │
       └──────────────────│          │ │ PDF/... │ │ 表达矩阵     │
                          └──────────┘ └────────┘ └────────────┘

核心原则不变:LLM 负责"决策",工具负责"执行"。区别是工具层现在对接多种数据源,每种源都有独立的白名单与校验,模型只决定"调哪个工具、传什么参数",永远不直接拼 SQL / 不碰原始文件路径拼装。

3. 技术选型

组件 选择 用途
语言 Python 3.11+ AI / 数据处理生态最成熟
LLM 接入 OpenAI 兼容 API(function calling + embeddings) 决策 + 文档语义检索
数据库 mysql-connector-python(只读账号) 结构化查询
表格 pandas Excel/CSV 解析与聚合
文档 pypdf / python-docx + 向量检索 文档抽取与语义搜索
生信文件 bcftools/samtools(CLI 包装)+ 纯 Python 解析 变异/序列/表达矩阵
编排 自写 Agent 循环 透明、可控、易调试

同样刻意不引入重型框架,用清晰的代码把"多源接入"讲透。

4. 统一的数据接入层

所有连接器的设计准则是:把"能取什么"固化成受控工具,路径与权限都在配置里白名单,模型只传语义化参数

4.1 数据库(只读 + 白名单)

# connectors/db.py
import os, re, mysql.connector

ALLOWED_TABLES = {"users", "orders", "posts"}
MAX_ROWS = 5000
conn = mysql.connector.connect(
    host=os.getenv("DB_HOST"),
    user=os.getenv("DB_USER_READONLY"),   # 只读账号
    password=os.getenv("DB_PASSWORD"),
    database=os.getenv("DB_NAME"),
)

def safe_query(table, columns, where=None, limit=100):
    if table not in ALLOWED_TABLES:
        raise ValueError(f"table not allowed: {table}")
    if where and not re.fullmatch(r"[\w ,=<>!'.%()]*", where):
        raise ValueError("invalid where clause")
    cols = ", ".join(c for c in columns if re.fullmatch(r"[\w]+", c))
    sql = f"SELECT {cols} FROM {table}"
    if where:
        sql += f" WHERE {where}"
    sql += f" LIMIT {min(limit, MAX_ROWS)}"
    cur = conn.cursor(dictionary=True)
    cur.execute(sql)
    rows = cur.fetchall()
    cur.close()
    return rows

4.2 表格数据(pandas,受控过滤)

表格工具用 pandas,但不用 df.query 直接执行模型字符串(有代码执行风险),而是实现一套安全过滤解析器,仅支持 列 运算符 值 的简单表达式。

# connectors/sheet.py
import os, re as _re, pandas as pd

SHEET_DIR = os.getenv("SHEET_DIR")           # 仅允许此目录下的文件
MAX_ROWS = 5000
ALLOWED_OPS = {"==", "!=", ">", "<", ">=", "<=", "contains"}
ALLOWED_AGGS = {"sum", "mean", "count", "max", "min", "median"}

def _resolve(name):
    path = os.path.abspath(os.path.join(SHEET_DIR, name))
    if not path.startswith(os.path.abspath(SHEET_DIR)):
        raise ValueError("path not allowed")
    return path

def _safe_filter(df, expr):
    # 顶层按 OR 拆分,每个 OR 块内按 AND 组合;只认 "列 运算符 值" 形式
    total_mask = pd.Series(False, index=df.index)
    for group in _re.split(r"\s+OR\s+", expr, flags=_re.I):
        terms = _re.split(r"\s+AND\s+", group, flags=_re.I)
        mask = pd.Series(True, index=df.index)
        for t in terms:
            m = _re.match(r"(\w+)\s*(==|!=|>=|<=|>|<|contains)\s*(.+)", t.strip())
            if not m:
                raise ValueError(f"invalid filter term: {t}")
            col, op, val = m.groups()
            val = val.strip().strip("'\"")
            if col not in df.columns:
                raise ValueError(f"unknown column: {col}")
            series = df[col]
            if op == "contains":
                sub = series.astype(str).str.contains(val, na=False)
            elif op in (">", "<", ">=", "<="):
                num = float(val)
                sub = (series > num) if op == ">" else (series < num) if op == "<" \
                    else (series >= num) if op == ">=" else (series <= num)
            else:
                sub = (series == val) if op == "==" else (series != val)
            mask &= sub
        total_mask |= mask
    return df[total_mask]

def query_sheet(filename, sheet=0, columns=None, where=None, group_by=None,
                agg=None, limit=100):
    path = _resolve(filename)
    df = pd.read_excel(path, sheet_name=sheet) if path.endswith((".xlsx", ".xls")) \
        else pd.read_csv(path)
    if columns:
        df = df[[c for c in columns if c in df.columns]]
    if where:
        df = _safe_filter(df, where)
    if group_by and agg:
        if agg not in ALLOWED_AGGS:
            raise ValueError(f"agg not allowed: {agg}")
        if group_by not in df.columns:
            raise ValueError(f"group_by not in columns: {group_by}")
        numeric = df.select_dtypes(include="number").columns.tolist()
        df = df.groupby(group_by)[numeric].agg(agg).reset_index()  # 白名单函数,作用于数值列
    return df.head(min(limit, MAX_ROWS)).to_dict(orient="records")

注:agg 已收敛为白名单函数名(sum/mean/count/max/min/median),作用于数值列;where 仅支持 列 运算符 值,顶层 OR、块内 AND,列名与运算符都做了校验,杜绝任意表达式与 eval

4.3 文档(抽取 + 向量检索)

文档不适合"逐字查询",更适合语义检索:把文档切成片段、做 embedding、按相似度召回最相关片段交给 LLM。

# connectors/doc.py
import os, numpy as np
from openai import OpenAI
import pypdf

DOC_DIR = os.getenv("DOC_DIR")
client = OpenAI(api_key=os.getenv("LLM_API_KEY"), base_url=os.getenv("LLM_BASE_URL"))

_chunks = []   # [(text, embedding)]

def _extract(path):
    if path.endswith(".pdf"):
        r = pypdf.PdfReader(path)
        return "\n".join(p.extract_text() or "" for p in r.pages)
    if path.endswith(".docx"):
        import docx
        d = docx.Document(path)
        return "\n".join(p.text for p in d.paragraphs)
    if path.endswith(".md"):
        return open(path, encoding="utf-8").read()
    raise ValueError("unsupported doc type")

def _chunk(text, size=800, overlap=100):
    return [text[i:i+size] for i in range(0, len(text), size-overlap)]

def build_index():
    """启动时一次性建索引(生产可持久化到向量库)。"""
    global _chunks
    for f in os.listdir(DOC_DIR):
        txt = _extract(os.path.join(DOC_DIR, f))
        for c in _chunk(txt):
            emb = client.embeddings.create(model="text-embedding-3-small",
                                           input=c).data[0].embedding
            _chunks.append((c, emb))

def search_documents(query, top_k=3):
    q = np.array(client.embeddings.create(model="text-embedding-3-small",
                                          input=query).data[0].embedding)
    scored = [(c, np.dot(q, e) / (np.linalg.norm(q) * np.linalg.norm(e)))
              for c, e in _chunks]
    scored.sort(key=lambda x: x[1], reverse=True)
    return [c for c, _ in scored[:top_k]]

4.4 系统特有文件 —— 以生信数据为例

生信文件形态特殊,且常常很大(BAM 可达数十 GB)。接入原则是:小文件纯 Python 解析;大文件(VCF.gz / BAM)包装 bcftools/samtools 命令行工具,不把整个文件读进内存。

# connectors/bio.py
import os, subprocess, re, csv

BIO_DIR = os.getenv("BIO_DIR")

def _resolve(name):
    path = os.path.abspath(os.path.join(BIO_DIR, name))
    if not path.startswith(os.path.abspath(BIO_DIR)):
        raise ValueError("path not allowed")
    return path

def vcf_summary(name):
    """变异统计:优先用 bcftools(支持大文件与压缩格式)。"""
    path = _resolve(name)
    out = subprocess.run(["bcftools", "view", "-H", path],
                         capture_output=True, text=True)
    types, samples = {}, set()
    for ln in out.stdout.splitlines():
        f = ln.split("\t")
        if len(f) < 8:
            continue
        info = f[7]
        m = re.search(r"TYPE=([^;]+)", info)
        t = m.group(1) if m else "UNKNOWN"
        types[t] = types.get(t, 0) + 1
        if len(f) > 9:
            samples.update(f[9:])
    return {"variant_count": len(out.stdout.splitlines()),
            "by_type": types, "sample_count": len(samples)}

def fasta_stats(name):
    """序列统计:纯 Python 流式读取,避免大文件占内存。"""
    path = _resolve(name)
    count, lengths, cur = 0, [], 0
    with open(path) as fh:
        for line in fh:
            if line.startswith(">"):
                if cur:
                    lengths.append(cur)
                count += 1
                cur = 0
            else:
                cur += len(line.strip())
        if cur:
            lengths.append(cur)
    return {"seq_count": count,
            "total_bp": sum(lengths),
            "max_len": max(lengths) if lengths else 0}

def expression_top(name, top_n=10):
    """表达矩阵 Top 基因:TSV,首列为 gene_id,其余为样本。"""
    path = _resolve(name)
    top = []
    with open(path) as fh:
        r = csv.reader(fh, delimiter="\t")
        header = next(r)
        for row in r:
            vals = [float(x) for x in row[1:]]
            top.append((row[0], sum(vals) / len(vals)))
    top.sort(key=lambda x: x[1], reverse=True)
    return [{"gene": g, "mean_expr": round(v, 3)} for g, v in top[:top_n]]

生信场景真实坑点:VCF/FASTA 常以 .gz 压缩、BAM 是二进制——别试图用 open() 直接读。能调 bcftools/samtools 就调,它们比手写的解析器更快也更准。

5. 工具定义与分发

把每类能力以 function calling schema 暴露。description 越清晰,模型选得越准

# tools.py
TOOLS = [
  {"type":"function","function":{"name":"query_table",
    "description":"查询网站业务库(users/orders/posts)的结构化数据,用于统计分析。",
    "parameters":{"type":"object","properties":{
      "table":{"type":"string","enum":["users","orders","posts"]},
      "columns":{"type":"array","items":{"type":"string"}},
      "where":{"type":"string","description":"如 \"status=1 AND created_at>='2026-08-01'\""},
      "limit":{"type":"integer"}},"required":["table","columns"]}}},
  {"type":"function","function":{"name":"query_sheet",
    "description":"查询后台表格文件(Excel/CSV),支持过滤与分组聚合。",
    "parameters":{"type":"object","properties":{
      "filename":{"type":"string"},
      "columns":{"type":"array","items":{"type":"string"}},
      "where":{"type":"string","description":"如 \"amount > 100 AND channel == 'wechat'\""},
      "group_by":{"type":"string","description":"分组列名"},
      "agg":{"type":"string","description":"聚合函数名:sum/mean/count/max/min/median,作用于数值列"},
      "limit":{"type":"integer"}},"required":["filename"]}}},
  {"type":"function","function":{"name":"search_documents",
    "description":"在政策/产品文档中语义检索相关片段。用于回答依据类问题。",
    "parameters":{"type":"object","properties":{
      "query":{"type":"string"},"top_k":{"type":"integer"}},"required":["query"]}}},
  {"type":"function","function":{"name":"vcf_summary",
    "description":"统计生信 VCF 变异文件的变异数、类型分布与样本数。",
    "parameters":{"type":"object","properties":{"name":{"type":"string"}},"required":["name"]}}},
  {"type":"function","function":{"name":"fasta_stats",
    "description":"统计 FASTA 序列文件的序列数与总碱基数。",
    "parameters":{"type":"object","properties":{"name":{"type":"string"}},"required":["name"]}}},
  {"type":"function","function":{"name":"expression_top",
    "description":"返回表达矩阵中平均表达量最高的若干基因。",
    "parameters":{"type":"object","properties":{
      "name":{"type":"string"},"top_n":{"type":"integer"}},"required":["name"]}}},
]

def dispatch(name, args):
    from connectors import db, sheet, doc, bio
    if name == "query_table":      return {"rows": db.safe_query(**args)}
    if name == "query_sheet":      return {"rows": sheet.query_sheet(**args)}
    if name == "search_documents": return {"chunks": doc.search_documents(**args)}
    if name == "vcf_summary":      return bio.vcf_summary(**args)
    if name == "fasta_stats":      return bio.fasta_stats(**args)
    if name == "expression_top":   return {"top_genes": bio.expression_top(**args)}
    return {"error": f"unknown tool: {name}"}

6. Agent 核心循环

与单源版本完全一致——模型决策 → 执行工具 → 结果回灌 → 再决策,直到收敛。多源只是工具更多。

# agent.py
import json, os
from openai import OpenAI
from tools import TOOLS, dispatch

client = OpenAI(api_key=os.getenv("LLM_API_KEY"), base_url=os.getenv("LLM_BASE_URL"))
MAX_TURNS = 10  # 防止模型陷入无限工具调用

SYSTEM_PROMPT = """你是网站数据分析助手,可访问数据库、表格、文档、生信文件等多源数据。
规则:
1. 优先用工具获取真实数据,禁止编造数字。
2. 跨源问题时,分别取数后再交叉分析,并说明数据来源。
3. 给出具体数字与可执行建议;必要时指出数据局限。
4. 中文回答。"""

def run_agent(question: str) -> str:
    messages = [{"role":"system","content":SYSTEM_PROMPT},
                {"role":"user","content":question}]
    for _ in range(MAX_TURNS):
        resp = client.chat.completions.create(
            model=os.getenv("LLM_MODEL","gpt-4o-mini"),
            messages=messages, tools=TOOLS, tool_choice="auto")
        msg = resp.choices[0].message
        if not msg.tool_calls:
            return msg.content
        messages.append(msg)
        for call in msg.tool_calls:
            args = json.loads(call.function.arguments)
            result = dispatch(call.function.name, args)
            messages.append({"role":"tool","tool_call_id":call.id,
                             "content":json.dumps(result, ensure_ascii=False, default=str)})
    return "(已达到最大推理轮数,请拆分问题后重试)"

7. 一个跨源真实例子

if __name__ == "__main__":
    q = ("上季度 A 类商品退货率偏高,请结合 sales_2026Q3.xlsx 的销售数据、"
         "数据库 orders 表、以及《退款政策.docx》,分析退货集中在哪些渠道;"
         "再调出生信样本 sample1.vcf.gz 看是否有相关变异标记。")
    print(run_agent(q))

模型可能的执行路径:

  1. query_sheet("sales_2026Q3.xlsx", where="category == 'A'") → 拿到 A 类销售与退货列;
  2. query_table("orders", ["channel","status"], where="status=3") → 拿到退货订单渠道分布;
  3. search_documents("退货 退款 政策 渠道") → 召回政策中关于各渠道退货的条款;
  4. vcf_summary("sample1.vcf.gz") → 拿到变异类型与样本信息;
  5. 综合四源数据,输出:

结论:A 类商品退货 68% 集中在"微信投放"渠道,与《退款政策》第 4 条"投放渠道 7 天无理由"高度相关;建议对该渠道加首单质检。生信样本 sample1 检出 3 个 high-impact 变异(TYPE=missense_variant),虽非直接退货原因,但可作为"样本质量"标记归档。 数据来源:sales_2026Q3.xlsx、orders 表、退款政策.docx、sample1.vcf.gz。

注意:每个数字都来自真实工具返回,文档片段与变异统计也都有据可查。

8. 进阶

  • 流式输出stream=True 逐字渲染。
  • FastAPI 接口:把 run_agent 包成 /api/agent/ask,前端聊天框直接调用,成为后台里的"AI 分析助手"。
  • 索引持久化:文档 embedding 索引落盘到向量库(FAISS / Milvus),避免每次启动重建。
  • 生信专用工具:把常用分析(如"按影响程度筛变异")固化成语义化工具 vcf_filter_by_impact(level),模型只传参数,更安全可解释。
  • 结果缓存:相同口径查询缓存到 Redis,降低 LLM 与 IO 压力。
  • 结果规模控制:工具可能返回上万行,直接喂给模型会撑爆上下文;后续(四)《记忆与上下文管理》会讲如何压缩与截断大结果。

9. 安全清单(多源尤其重要)

  • 数据库用只读账号 + 表白名单 + where 正则校验
  • 表格:文件路径白名单(限定 SHEET_DIR),过滤表达式走安全解析器,禁用任意 eval
  • 文档:限定 DOC_DIR,抽取与检索分离,embedding key 不进前端
  • 生信文件:限定 BIO_DIR;大文件一律走 bcftools/samtools,禁止全量读内存
  • 所有工具返回做行数 / 大小上限,防止资源耗尽
  • 审计日志:记录每次提问、调用的工具、参数与返回摘要
  • 对外暴露时加鉴权,避免成为数据泄露口

10. 小结

多源 Agent 的难度不在"模型多聪明",而在把每类数据的访问收敛成受控、白名单化的工具。只要守住"模型只决策、工具才执行、每源独立护栏"这条线,就能让 AI 安全地坐在你的数据库、表格、文档乃至生信文件之上,把零散的跨源分析需求变成一句自然语言。

可继续演进:把文档索引迁入向量库、为生信场景增加更多专用工具、用 RAG 让文档回答更精准、把结论自动生成日报推送。但无论怎么扩展,先界定"Agent 能碰什么",再谈"能回答什么"

本篇在(一)的最小骨架上接入了真实数据,是系列里"接数据"这一层;后面每篇都在它之上叠加一层能力——(三)工具编排、(四)记忆、(五)前端、(六)检索、(七)多 Agent、(八)观测、(九)部署——它们共享同一套"模型决策、工具执行、每源护栏"的底座。

给网站加 AI,最容易翻车的不是模型效果,而是数据权限与多源边界。护栏先行,跨源才稳。

下一篇(三)《工具系统深入》:我们会把本篇的 dispatch 从"能调"升级到"稳、快、可控"——并行调用、失败重试、权限分级、结果缓存与超时控制。