大山佬头像
关注

Data Agent 架构:让 LLM 自主完成取数与分析的闭环

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

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--