Agent 开发实战(二):对接多源后台数据的分析 Agent
Agent 开发实战(二):对接多源后台数据的分析 Agent
真实网站的"后台数据"从来不止一张数据库表。它可能是 MySQL 里的业务表、Excel/CSV 报表、产品/政策文档(PDF/Word/Markdown),也可能是某个垂直系统特有的文件——比如生信分析系统导出的 VCF 变异文件、FASTA 序列、表达矩阵等。本文手把手实现一个统一接入多源数据、用 LLM 做跨源分析的 Agent,并重点讲清楚每类数据源的安全接入方式。
这是「Agent 开发实战」系列的第 2 篇,承接(一)跑通的最小 Agent。(一)已经把"模型决策 / 工具执行"的核心循环讲清楚了,本篇直接复用这个模式,重点放在把多种数据源变成受控工具、让 Agent 能回答跨源问题。后续(三)~(九)会依次深入工具编排、记忆、前端、RAG、多 Agent、可观测性与工程化部署。
系列地图:(一)最小 Agent | (二)多源数据接入 ← 本篇 | (三)工具系统深入 | (四)记忆与上下文 | (五)流式与前端集成 | (六)RAG 进阶 | (七)多 Agent 协作 | (八)可观测性与评估 | (九)工程化部署与端到端案例。逻辑是「地基 → 接数据 → 工具工程化 → 记忆 → 前端 → 检索 → 多 Agent → 观测 → 部署」,每篇都基于前一篇叠加一层能力,但共享同一套"模型决策、工具执行、每源护栏"的底座。
1. 我们要解决什么
一个典型网站后台,数据散落在不同形态里:
- 数据库:
users、orders、posts等结构化业务表 - 表格数据:运营导出的
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))
模型可能的执行路径:
query_sheet("sales_2026Q3.xlsx", where="category == 'A'")→ 拿到 A 类销售与退货列;query_table("orders", ["channel","status"], where="status=3")→ 拿到退货订单渠道分布;search_documents("退货 退款 政策 渠道")→ 召回政策中关于各渠道退货的条款;vcf_summary("sample1.vcf.gz")→ 拿到变异类型与样本信息;- 综合四源数据,输出:
结论: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从"能调"升级到"稳、快、可控"——并行调用、失败重试、权限分级、结果缓存与超时控制。