基于Agent框架构建AI数据医生:实现数据平台智能运维闭环
1. 从“救火队员”到“数据医生”的诞生
去年年底,我们团队的数据仓库又双叒叕出问题了。凌晨两点,我被一连串的告警电话叫醒,业务报表延迟、ETL任务大面积失败、查询响应时间飙升到分钟级。等我连上服务器,面对满屏的日志和监控图表,花了整整三个小时,才定位到是一个上游数据源的字段类型变更,导致下游十几个依赖它的任务链式崩溃。这已经不是第一次了,团队里每个数据工程师都扮演过“救火队员”的角色,疲于奔命地处理各种数据质量问题:字段值异常、数据延迟、任务失败、存储空间告警……
痛定思痛,我意识到,靠人力去监控和修复海量、复杂的数据管道,效率低下且不可持续。我们需要一个能7x24小时“值班”的智能体,它不仅能像医生一样“望闻问切”——实时监测数据健康状况,更要能“开方抓药”——自动诊断问题根因并执行修复动作。这就是我们内部称之为“AI数据医生”项目的起源。它不是某个单一的工具,而是一个基于Agent(智能体)框架构建的、具备感知、诊断、决策和执行能力的自动化运维系统。今天,我就来拆解一下这个“数据医生”的核心架构和关键配置,尤其是我们如何利用OpenClaw这类框架来构建它的“大脑”和“手脚”。整个过程,充满了从零到一的探索和踩坑,希望我们的经验能帮你少走弯路。
2. “数据医生”的顶层设计:一个具备闭环能力的智能体
在开始敲代码之前,我们必须想清楚这个“医生”到底要干什么。它不是一个简单的监控告警系统,后者只负责“发现问题并通知人”。我们的目标是让系统能“发现问题并尝试解决问题”,如果解决不了,再清晰地告诉人“病根”在哪里以及它已经做了哪些尝试。这本质上是一个感知-思考-行动的循环,也就是现在常说的Agent(智能体)架构。
2.1 核心能力定义
我们的“数据医生”被赋予了四大核心能力:
感知与检查(Examination):这是医生的“望闻问切”。系统需要持续从各个“器官”(数据源)收集“生命体征”。这包括:
- 批处理任务:Airflow/Dagster任务的状态、日志、运行时长。
- 流处理管道:Flink/Kafka Consumer的Lag、吞吐量、错误率。
- 数据存储:Hive/MySQL/ClickHouse的表大小增长趋势、分区健康度、主键重复率。
- 数据质量:字段的空值率、值域分布、与历史同期的统计差异(如平均值突增)、数据新鲜度(最后更新时间)。
- 计算资源:YARN/K8s队列的资源使用率、Spark任务的GC情况。
诊断与推理(Diagnosis):这是医生的“病情分析”。收集到异常指标后,需要判断这是孤立事件还是系统性问题,并找到最可能的根因。例如,“Hive表A的写入任务失败”可能的原因有:上游数据格式错误、HDFS空间不足、Hive Metastore连接超时、计算资源不足等。单纯的规则引擎(如果失败则告警)在这里不够用,我们需要一个能联系上下文进行推理的“大脑”。
决策与规划(Planning):确诊后,需要决定“治疗方案”。是自动重试任务?是执行一个数据修复脚本?是扩容计算资源?还是需要人工介入?决策需要基于预设的策略和成本评估。例如,对于非核心表的短暂延迟,可以自动重试;对于核心财务数据的错误,则必须立即告警并阻止下游消费。
执行与反馈(Action & Feedback):开出“药方”后,要能“抓药”。系统需要安全地执行各种操作,如调用K8s API重启Pod、在Airflow中触发任务重跑、执行一段SQL修复数据、或发送一条结构清晰的告警信息到飞书/钉钉。执行后,还要再次“检查”,确认治疗是否有效,形成闭环。
2.2 技术架构选型:为什么是Agent框架?
要实现上述能力,我们有几种选择:
- 纯脚本调度:写一堆Crontab和Shell脚本。缺点:逻辑僵化,难以维护和扩展,缺乏推理能力。
- 工作流引擎(如Airflow):擅长编排确定性的任务流,但对异常处理和非确定性决策支持弱。
- 规则引擎(如Drools):可以处理“如果-那么”规则,但面对复杂、多变的上下文,规则库会爆炸,且难以维护。
最终我们选择了Agent框架。因为它天然契合“感知-思考-行动”范式。一个Agent封装了状态、策略和能力,能在环境中通过工具(Tools)获取信息(感知),根据目标进行推理(思考),并选择工具执行动作(行动)。近年来兴起的大语言模型(LLM)为Agent提供了强大的、泛化的推理能力,使其能够理解自然语言描述的问题,并规划出解决步骤。
在我们的场景中,“数据医生”就是一个LLM驱动的自主Agent。LLM作为其“大脑”,负责复杂的诊断和规划;而一系列封装好的数据平台API和运维脚本则作为其“工具”,供大脑调用以执行具体操作。
3. 构建“大脑”:基于OpenClaw框架的智能体核心
框架选型上,我们评估了LangChain、LlamaIndex、AutoGen以及OpenClaw。OpenClaw吸引我们的点在于它对生产级AI应用的开箱即用支持,特别是其清晰的Agent-Skill-Tool分层架构和面向企业集成的设计。
注意:框架选型没有绝对优劣,取决于团队技术栈和场景。LangChain生态最繁荣但略显臃肿;LlamaIndex对数据索引专注;AutoGen多Agent对话强大。OpenClaw的“技能”抽象与我们“数据医生”的“专科能力”想法不谋而合。
3.1 OpenClaw核心概念与我们的映射
- Agent(智能体):这就是我们的“数据医生”本体。在OpenClaw中,一个Agent由配置定义,包括其名称、描述、使用的LLM、拥有的技能(Skills)以及记忆(Memory)等。我们创建了一个名为
DataPlatformDoctor的Agent。 - Skill(技能):这是Agent的“专科”。我们将“数据医生”的四大能力拆解成不同的Skill:
DataHealthExaminationSkill:感知与检查技能。负责调用各类监控API获取数据。AnomalyDiagnosisSkill:异常诊断技能。封装诊断逻辑,可以调用规则引擎,也可以让LLM分析。RemediationPlanningSkill:修复规划技能。基于诊断结果,制定行动计划。PlatformActionSkill:平台操作技能。负责安全地执行各类运维动作。
- Tool(工具):Skill的具体实现手段。一个Skill可以调用多个Tool。例如,
DataHealthExaminationSkill可能包含以下Tool:query_prometheus_tool: 查询Prometheus获取指标。fetch_airflow_dag_status_tool: 获取Airflow DAG运行状态。check_hive_table_metadata_tool: 检查Hive表元数据健康度。
- LLM(大语言模型):Agent的推理引擎。我们通过OpenClaw的配置,让Agent使用部署在内部的ChatGLM3或通义千问的API。对于诊断和规划这类需要复杂推理的环节,LLM是核心。
3.2 核心配置文件详解
OpenClaw的配置是其强大之处。以下是我们DataPlatformDoctor的核心配置片段,我加了详细注释:
# config/agent_data_doctor.yaml agent: name: "DataPlatformDoctor" description: "一个自动监控、诊断和修复数据平台问题的AI智能体。" # 使用本地部署的LLM,避免网络延迟和依赖 llm: type: "openai" # OpenClaw兼容OpenAI API协议 base_url: "http://your-llm-api-server/v1" # 内部ChatGLM/Qwen API地址 model: "qwen-plus" # 模型名称 api_key: "${LLM_API_KEY}" # 从环境变量读取密钥 # 技能列表:这是“数据医生”的专科能力 skills: - "skill_data_health_examination" - "skill_anomaly_diagnosis" - "skill_remediation_planning" - "skill_platform_action" # 记忆:让Agent能记住近期处理过的问题,避免重复动作 memory: type: "short_term" window_size: 10 # 记住最近10轮对话/事件 # 技能的具体配置 skill_data_health_examination: description: "从各类监控系统收集数据平台健康指标。" tools: - "tool_prometheus_query" - "tool_airflow_api" - "tool_hive_metastore_client" # 触发条件:每5分钟自动运行一次,也可由告警事件触发 triggers: - type: "cron" expression: "*/5 * * * *" - type: "event" pattern: "alert.fired" skill_anomaly_diagnosis: description: "分析健康指标,识别异常并推断根因。" # 这个技能严重依赖LLM进行推理 llm_usage: "high" tools: - "tool_rule_engine" # 先过一遍硬规则,过滤掉明显问题 - "tool_llm_analyzer" # 复杂情况交给LLM分析 # 配置诊断提示词模板,这是让LLM当好“医生”的关键 prompt_templates: diagnosis: | 你是一个资深数据平台运维专家。请分析以下数据平台的异常情况,并给出最可能的根本原因。 上下文信息: - 异常指标: {metrics} - 近期事件: {recent_events} - 相关任务/表: {affected_entities} 请按以下格式回答: 1. 根因分析:[你的推理过程] 2. 置信度:[高/中/低] 3. 建议的下一步行动:[例如:检查X日志、执行Y操作、通知Z人员]这个配置定义了Agent的骨架。其中,提示词模板(Prompt Template)是灵魂。我们花了大量时间迭代diagnosis模板,通过提供结构化的上下文(指标、事件、实体)和要求结构化的输出,极大地提升了LLM诊断的准确性和可用性。
3.3 工具(Tool)的实现与安全考量
工具是Agent与真实世界交互的手脚。实现它们时,安全性是第一位。我们绝不允许Agent拥有不受限制的root权限。
以tool_airflow_api为例,我们不是直接让Agent持有Airflow的admin密码,而是实现了一个轻量的代理服务:
# tools/airflow_tool.py import requests from openclaw.tool import tool @tool def trigger_dag_run(dag_id: str, conf: dict = None) -> dict: """ 触发一个Airflow DAG运行。 安全策略:仅允许触发标签为‘auto_remediable’的DAG。 """ # 1. 安全检查:查询DAG信息,确认其标签 dag_info = _get_dag_info(dag_id) if 'auto_remediable' not in dag_info.get('tags', []): return {"error": f"DAG {dag_id} 未标记为可自动修复,操作被拒绝。"} # 2. 通过一个具有严格权限的服务账号调用Airflow REST API url = f"{AIRFLOW_WEB_URL}/api/v1/dags/{dag_id}/dagRuns" payload = {"conf": conf} if conf else {} headers = {"Authorization": f"Bearer {AIRFLOW_SERVICE_ACCOUNT_TOKEN}"} response = requests.post(url, json=payload, headers=headers) return response.json() def _get_dag_info(dag_id: str) -> dict: # 内部方法,用于获取DAG元数据 pass所有具备“写”能力的工具(如重启服务、修复数据)都必须经过类似的二次校验。我们为“数据医生”创建了专属的、权限最小化的服务账号,并在工具层实现业务逻辑层面的安全策略(如只允许重试特定标签的任务)。
4. “诊断”流程的工程化实现:从规则到推理
单纯的LLM调用不稳定,且成本高。我们的诊断流程是规则引擎先行,LLM兜底的混合模式。
4.1 第一层:基于规则的快速过滤
我们使用一个轻量级规则引擎(比如自研的简单规则匹配,或接入开源方案如Drools)处理那些模式明确的常见问题。规则用YAML配置,易于运维同学管理。
# rules/common_failures.yaml rules: - name: "hdfs_disk_full_etl_fail" condition: | metrics.hdfs_disk_usage_percentage > 90 && events.last_hour contains “ETL job failed with IOException” diagnosis: “HDFS存储空间不足导致ETL任务写入失败。” suggested_action: “清理过期数据或扩容HDFS存储。” priority: “HIGH” auto_action: “trigger_cleanup_dag” # 关联到自动修复的DAG - name: “kafka_consumer_lag_spike” condition: | metrics.kafka_consumer_lag{consumer_group=“flink_job”} > 10000 && metrics.flink_task_manager_cpu > 85 diagnosis: “Flink消费者处理速度下降,可能由于资源不足或业务逻辑阻塞。” suggested_action: “检查Flink作业背压情况及TaskManager GC日志。” priority: “MEDIUM”规则引擎能毫秒级响应,解决80%的典型问题。如果所有规则都不匹配,或者规则匹配的置信度较低,事件才会被送入第二层。
4.2 第二层:LLM驱动的深度诊断
对于规则无法覆盖的复杂、新型问题,我们启动LLM诊断流程。这里的关键是为LLM构建高质量的上下文。我们不能只扔给它一条错误日志。
我们的tool_llm_analyzer工具会做以下工作:
- 信息聚合:从监控系统(Prometheus)、任务调度器(Airflow)、元数据库(Hive Metastore)等多个源头,收集与异常实体(如表、任务)相关的近期指标、日志片段、配置变更记录。
- 上下文构建:将上述信息整理成一段结构化的自然语言描述,填充到
skill_anomaly_diagnosis的提示词模板中。 - 调用与解析:调用配置的LLM,并解析其返回的结构化结果(根因分析、置信度、建议行动)。
这个过程中最大的坑是LLM的“幻觉”。它可能给出一个听起来合理但完全错误的诊断。我们的应对策略是:
- 要求结构化输出:如前文配置所示,强制要求LLM按固定格式回答,方便程序解析和校验。
- 设置置信度阈值:只有置信度为“高”的诊断,才会触发自动修复行动;“中”置信度的诊断会生成报告供人工审核;“低”置信度则仅记录,并补充更多监控信息。
- 人工反馈闭环:在飞书告警卡片上,我们增加了“诊断是否正确”的按钮。运维同学的反馈会被收集,用于后续优化提示词和规则。
5. “治疗”执行:安全、可控的自动化修复
诊断之后是治疗。我们严格遵循“可观测、可回滚、可干预”的原则来设计自动化修复动作。
5.1 修复动作的封装:标准化DAG与脚本
所有自动化修复操作,都不允许Agent直接执行Shell命令。而是封装成标准的、可重复执行的Airflow DAG或可版本化管理的脚本。
例如,对于“清理HDFS过期分区”这个动作,我们有一个预定义的DAGcleanup_old_hdfs_partitions。当规则引擎或LLM诊断建议此操作时,skill_platform_action技能只是去触发这个DAG的运行。DAG本身包含了完整的逻辑:确定哪些分区可删、创建备份(如需)、执行删除、验证清理结果。这样,修复逻辑是透明、可审查、可回滚的(备份)。
5.2 审批与熔断机制
并非所有修复都能全自动。我们在Agent的决策流程中加入了审批链概念。
# config/action_policies.yaml actions: - name: “retry_failed_task” auto_approval: true # 重试任务,低风险,自动执行 - name: “restart_flink_job” auto_approval: false approval_channel: “feishu_platform_ops” # 需要平台运维组在飞书审批 timeout: “5m” # 5分钟内无审批则升级告警 - name: “drop_production_table” auto_approval: false approval_channel: “feishu_data_owner_and_platform_lead” # 需数据负责人和平台负责人双重审批同时,我们为Agent设置了全局熔断器。如果短时间内(如10分钟内)自动修复失败次数超过阈值(如3次),Agent会自动进入“只监控,不行动”的安全模式,并发出最高级别告警,防止问题在自动化下被放大。
6. 部署与运维:让“数据医生”稳定服务
我们使用Docker容器化部署整个“数据医生”应用,核心组件包括:OpenClaw Agent服务、规则引擎服务、工具代理服务(封装各种平台API调用)。
6.1 部署架构
# docker-compose.prod.yaml 核心部分 version: ‘3.8’ services: openclaw-agent: image: our-registry/openclaw-agent:latest container_name:>