公司动态
Data Agent 架构:让 LLM 自主完成取数与分析的闭环
Data Agent 架构让 LLM 自主完成取数与分析的闭环一、取数分析的人肉中转困局数据团队每天都在重复一种低价值的劳动。业务方提出一个分析需求。工程师把它翻译成 SQL再等待查询结果。最后整理成表格或图表。这一过程里真正产生洞察的环节往往只占全部耗时的两成。绝大多数时间消耗在需求澄清、口径对齐与反复返工上。需求方并不懂底层的表结构和字段含义。工程师也不完全清楚业务方口中的活跃用户究竟指什么。双方靠自然语言来回拉扯。一个本该十分钟解决的问题常常拖成半天的扯皮。取数往往不是一步到位。业务方拿到第一张表后会顺着数据继续追问。为什么华东区下滑是不是某个渠道异常能否拆到城市粒度每一轮追问都要重新走一遍提需求—写 SQL—等结果的老路。这正是 Data Agent 要解决的命题。把 LLM 放在取数与分析的中枢位置。让它自己规划步骤、调用工具、读取结果、反思偏差直到产出可信的结论。工程师从人肉 SQL 翻译机退回到工具与口径的供给方与守门人。本文不讨论玩具级的 ChatGPT 套壳。重点放在可落地的生产架构如何让 LLM 在一个受控的闭环里安全、可观测、可回滚地完成多步分析。二、ReAct 闭环规划—执行—反思的自主取数引擎Data Agent 的核心是一个 ReActReason 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(okFalse, errorf工具 {name} 不在白名单内) if spec.risk RiskLevel.FORBIDDEN: return ToolResult(okFalse, 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), timeoutspec.timeout, ) if result is None: return ToolResult(okFalse, error工具返回空值视为失败) return ToolResult(okTrue, dataresult) 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(okFalse, errorlast_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(semaphore4) 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数据平台的杠杆率才能真正落到实处。