Data Agent 架构:让 LLM 自主完成取数与分析的闭环
一、取数分析的"人肉中转"困局
数据团队每天都在重复一种低价值的劳动。业务方提出一个分析需求。工程师把它翻译成 SQL,再等待查询结果。最后整理成表格或图表。这一过程里,真正产生洞察的环节,往往只占全部耗时的两成。
绝大多数时间,消耗在需求澄清、口径对齐与反复返工上。需求方并不懂底层的表结构和字段含义。工程师也不完全清楚业务方口中的"活跃用户"究竟指什么。双方靠自然语言来回拉扯。一个本该十分钟解决的问题,常常拖成半天的扯皮。
取数往往不是一步到位。业务方拿到第一张表后,会顺着数据继续追问。为什么华东区下滑?是不是某个渠道异常?能否拆到城市粒度?每一轮追问,都要重新走一遍"提需求—写 SQL—等结果"的老路。
这正是 Data Agent 要解决的命题。把 LLM 放在取数与分析的中枢位置。让它自己规划步骤、调用工具、读取结果、反思偏差,直到产出可信的结论。工程师从"人肉 SQL 翻译机",退回到"工具与口径的供给方与守门人"。
本文不讨论玩具级的 ChatGPT 套壳。重点放在可落地的生产架构:如何让 LLM 在一个受控的闭环里,安全、可观测、可回滚地完成多步分析。
二、ReAct 闭环:规划—执行—反思的自主取数引擎
Data Agent 的核心,是一个 ReAct(Reason + Act)循环。模型先推理当前状态,再决定调用哪个工具。工具返回结果后,模型反思是否足够,不够就继续规划下一步。
整个闭环由五个角色构成:规划器负责拆解目标;工具中枢负责执行白名单内的动作;记忆体保存中间结果;反思器评估置信度;守卫器拦截越权与高危操作。下面用一张时序图,呈现一次"华东区 GMV 为何下滑"的完整推理链。
sequenceDiagram
participant U as 用户提问
participant P as 规划器Planner
participant L as LLM推理核心
participant T as 工具中枢ToolHub
participant G as 守卫器Guard
participant M as 记忆体Memory
U->>L: 提交分析目标:华东区GMV下滑归因
L->>P: 拆解为多步子任务
P-->>L: 步骤1=取华东日级GMV
L->>G: 申请调用query_sql工具
G-->>L: 校验通过(命中白名单)
L->>T: 执行SQL查询
T-->>L: 返回近30日GMV序列
L->>M: 写入中间结论
M-->>L: 确认存储
L->>L: 反思:环比下降但渠道未拆分
L->>G: 申请调用query_sql(按渠道下钻)
G-->>L: 校验通过
T-->>L: 返回分渠道GMV
L->>M: 更新结论:社群渠道异常
L->>U: 输出归因报告与证据链
style P fill:#4A90D9,color:#fff
style L fill:#7B61FF,color:#fff
style T fill:#50C878,color:#fff
style G fill:#E0573E,color:#fff
style M fill:#F2B705,color:#000
工具白名单是闭环安全性的基石。Agent 永远不能直接触碰数据库 DDL。它能调用的,只有预先注册、参数受限、只读优先的若干工具。下表给出一组典型白名单与对应的风险等级。
flowchart LR
A[业务提问] --> B{工具中枢}
B -->|白名单内| C[query_sql只读查询]
B -->|白名单内| D[get_schema取表结构]
B -->|白名单内| E[plot_chart生成图表]
B -->|白名单内| F[calc_expr指标计算]
B -->|越权拦截| G[拒绝执行DDL/DML]
C --> H[结果回流记忆体]
D --> H
E --> H
F --> H
H --> I[反思器评估置信度]
style C fill:#50C878,color:#fff
style D fill:#50C878,color:#fff
style E fill:#50C878,color:#fff
style F fill:#50C878,color:#fff
style G fill:#E0573E,color:#fff
style I fill:#F2B705,color:#000
反射阶段必须设置"退出条件"。否则 Agent 会在不确定时无限循环。约定三类停止信号:置信度达标、步数触顶、或连续两轮结果无新增信息。
三、生产级 Data Agent 工具调度实现
下面给出工具中枢与守卫器的参考实现。代码覆盖超时控制、失败重试、并发上限、空值与异常兜底。实际工程中,query_sql 应指向受权限隔离的只读代理,而非直连生产库。
import asyncio
import time
from dataclasses import dataclass, field
from enum import Enum
from typing import Callable, Dict, Optional
import logging
logger = logging.getLogger("data_agent")
class RiskLevel(Enum):
READ_ONLY = "read_only" # 只读查询,风险最低
SCHEMA = "schema" # 读取元数据,风险低
COMPUTE = "compute" # 内存计算,风险低
FORBIDDEN = "forbidden" # 任何写操作,一律拦截
@dataclass
class ToolSpec:
name: str
handler: Callable
risk: RiskLevel
timeout: float = 8.0 # 单次调用超时秒数
max_retry: int = 2 # 失败重试上限
@dataclass
class ToolResult:
ok: bool
data: Optional[dict] = None
error: str = ""
class ToolHub:
"""受白名单约束的工具中枢,负责调度与限流。"""
def __init__(self, semaphore: int = 4):
self._tools: Dict[str, ToolSpec] = {}
# 全局并发信号量,避免瞬时打满查询引擎
self._sem = asyncio.Semaphore(semaphore)
def register(self, spec: ToolSpec) -> None:
if spec.risk == RiskLevel.FORBIDDEN:
raise ValueError(f"工具 {spec.name} 属于禁用类别,拒绝注册")
self._tools[spec.name] = spec
async def invoke(self, name: str, **kwargs) -> ToolResult:
spec = self._tools.get(name)
if spec is None:
return ToolResult(ok=False, error=f"工具 {name} 不在白名单内")
if spec.risk == RiskLevel.FORBIDDEN:
return ToolResult(ok=False, error="命中守卫器:拒绝执行高危操作")
async with self._sem:
return await self._run_with_retry(spec, **kwargs)
async def _run_with_retry(self, spec: ToolSpec, **kwargs) -> ToolResult:
last_err = ""
for attempt in range(spec.max_retry + 1):
try:
# 用 wait_for 实现硬超时,防止慢查询拖垮整个闭环
result = await asyncio.wait_for(
self._to_coro(spec.handler, **kwargs),
timeout=spec.timeout,
)
if result is None:
return ToolResult(ok=False, error="工具返回空值,视为失败")
return ToolResult(ok=True, data=result)
except asyncio.TimeoutError:
last_err = f"调用超时({spec.timeout}s),第{attempt + 1}次"
logger.warning(last_err)
except Exception as exc: # 兜底捕获,避免单工具崩溃中断 Agent
last_err = f"工具异常:{exc},第{attempt + 1}次"
logger.warning(last_err)
return ToolResult(ok=False, error=last_err)
@staticmethod
def _to_coro(handler: Callable, **kwargs):
if asyncio.iscoroutinefunction(handler):
return handler(**kwargs)
# 同步函数包成协程,统一异步调度
async def _wrap():
return handler(**kwargs)
return _wrap()
# 示例:只读 SQL 查询工具,需由上层注入受控连接
def make_query_sql(engine):
def _run(sql: str) -> dict:
# 真实环境应在此做 AST 校验,禁止 INSERT/UPDATE/DELETE
with engine.connect() as conn:
rows = conn.execute(text(sql)).mappings().all()
if not rows:
return {"rows": [], "count": 0}
return {"rows": [dict(r) for r in rows], "count": len(rows)}
return _run
if __name__ == "__main__":
hub = ToolHub(semaphore=4)
hub.register(ToolSpec("query_sql", make_query_sql(None), RiskLevel.READ_ONLY))
hub.register(ToolSpec("get_schema", lambda t: {"table": t}, RiskLevel.SCHEMA))
记忆体建议用带 TTL 的键值存储。每轮反思把"已确认事实"与"待验证假设"分桶保存。这样即便中途失败,也能从最近的检查点恢复,而非从头再来。
四、边界条件、Trade-offs 与适用禁用
任何架构都有它的适用面。Data Agent 同样如此。
边界条件方面,首当其冲的是口径幻觉。LLM 可能把"活跃用户"理解成"登录过即算",而企业口径是"完成关键行为"。这种偏差不会报错,却会 quietly 误导决策。必须用指标字典作为硬约束注入提示词。其次是长链路失控:当拆解超过十步,反思器容易陷入自我否定循环。需要设置最大步数硬上限。
Trade-offs 需要清醒权衡。引入 Agent 后,单次分析延迟通常高于人工直写 SQL,因为多了推理与重试开销。换来的是需求方的自助能力与响应速度。在查询成本上,Agent 的试探性查询会产生冗余扫描,需要靠结果缓存与查询去重来压低开销。
适用场景包括:口径相对固化、表结构稳定、以探索式分析为主的业务自助取数;临时性的多维下钻归因;以及需要把分析过程留痕、形成可复用报告模板的场合。
禁用场景必须明确划界。涉及财务结算、监管报送等强一致性要求的取数,不能交给概率模型自由发挥。涉及未脱敏 PII 的明细查询,工具层必须事先拦截。涉及 DDL/DML 的库表变更,无论如何都不应进入白名单。
一个务实的落地节奏是:先让 Agent 处理"只读、可解释、可回放"的轻量分析,把高风险动作保留给人工。待评估体系成熟,再逐步放开边界。
五、总结
Data Agent 不是要取代数据工程师。它的价值,在于把重复、低价值、易扯皮的取数环节自动化。让专业的人,回到真正需要判断力的地方。
落地的关键,不在模型有多大,而在闭环是否受控。白名单约束工具边界,反思器控制推理深度,守卫器兜住安全底线。三者齐备,ReAct 闭环才从演示走向生产。
架构图先行、代码紧随、原理收尾。当取数与分析的闭环被稳妥地交给 Agent,数据平台的杠杆率才能真正落到实处。
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/qq_42431428/article/details/163399626



