NANDAN-CREATOR/pipeline-incident-agent
GitHub: NANDAN-CREATOR/pipeline-incident-agent
一个基于本地 LLM 的数据管道事件响应 Agentic AI 示例,演示如何自动诊断管道失败并提供带证据链的修复建议。
Stars: 0 | Forks: 0
# Pipeline 事件响应 Agent(示例 / Demo)
这是一个应用于数据工程问题的**Agentic AI 循环**的最小化端到端演示:像值班工程师一样诊断失败的数据管道任务——检查日志、提出假设、从多个来源收集佐证证据,最终提出修复建议或上报给人工处理,且全程不会在无人监督的情况下执行任何操作。
它完全通过 **[Ollama](https://ollama.com) 运行在本地模型**上——无需云 API 密钥。
## 演示内容
一个夜间 pipeline 的 join 步骤失败并出现以下错误:
```
KeyError: 'customer_region'
```
Agent 不会等待人工立即深入排查日志,而是自行展开调查:
1. 获取失败任务的日志并发现 `KeyError`
2. 形成假设:schema 发生变更,或者上游任务静默失败
3. 检查相关表的 schema 变更历史
4. **不满足于第一个看似合理的答案**——还会检查上游任务状态,以排除其他可能的解释
5. 对线上真实数据执行只读查询,以确认 schema 历史与实际情况相符(因为元数据可能会滞后于真实数据)
6. 检查近期的代码变更,确认是否已经有人进行了修复
7. 提出具体且最小化的修复方案——包含置信度和完整的证据链——供**人工审查和批准**
Agent 绝不能自行应用修复、更改数据或重启任何内容。它最终只可能执行两种操作:`propose_fix` 和 `escalate_to_human`——这两者都仅仅是生成一条结构化记录,交由人工来采取行动。
## 环境要求
- Python 3.9+
- 已安装并在本地运行 [Ollama](https://ollama.com)
- 在 Ollama 中拉取一个具备工具调用能力的模型,例如:
ollama pull llama3.1
在本文撰写时,已知在 Ollama 中支持工具调用的其他模型包括 `qwen2.5`、`mistral-nemo` 和 `firefunction-v2`——在选择模型之前,请查看 [Ollama 模型库](https://ollama.com/library) 以获取当前的工具调用支持情况。较小的模型在持续生成格式良好的工具调用方面可能会明显不太可靠;如果您发现 Agent 立即上报并提示“未生成工具调用”,请先尝试使用更大或功能更强的模型,然后再排查其他问题。
## 安装设置
```
git clone https://github.com/NANDAN-CREATOR/pipeline-incident-agent.git
cd pipeline-incident-agent
pip install -r requirements.txt
# 在一个单独的终端中,确保 Ollama 正在运行:
ollama serve
# 确保你选择的模型已被拉取:
ollama pull llama3.1
```
## 运行 Demo
```
python agent.py
```
或者使用其他模型运行:
```
python agent.py --model qwen2.5
```
您将看到逐步的追踪信息打印在控制台上:模型在每一步决定调用哪个工具、该工具返回了什么结果,最后是 `propose_fix` 记录或 `escalate_to_human` 记录。
**关于输出的说明:** 由于模型在每一步都在做出真实的决策,因此每次运行及不同模型之间,具体的步骤数量和检查顺序可能会略有不同——这正是*Agentic*循环与固定脚本的区别所在。在合理的运行中,结论应当保持一致:将 `customer_region` 更改为 `region_code` 的高置信度修复,并得到至少两个相互印证的工具结果支持,因为这是系统提示词中的证据规则在调用 `propose_fix` 之前所必须满足的条件。
## 项目结构
```
pipeline-incident-agent/
├── agent.py # everything: mock tools, schemas, system prompt, agent loop
├── ARTICLE.md # the full companion article (design walkthrough)
├── requirements.txt
├── LICENSE
└── README.md
```
对于像这样的示例项目,我们特意将其简化为单个文件——关于在实际部署中如何对其进行拆分,请参阅[适配到真实系统](#adapting-this-to-real-systems)。
## 护栏机制如何运作
- **自由读取,审批后操作:** 除了两个终止类工具外,其他所有工具在设计上都是只读的。无论 Agent 做出什么决定,都没有任何工具可以用来修改任何内容。
- **证据纪律:** 系统提示词要求,在调用 `propose_fix` 之前,必须有至少两个独立的、相互印证的证据来源。仅仅有一个看似合理的原因是不够的。
- **提示词注入意识:** 系统提示词明确指示模型将工具输出视为数据,绝不能当作指令来执行——这一点非常关键,因为日志消息和其他工具输出正是真实系统在读取时需要格外小心的那种不受信任的内容。
- **步骤限制:** 如果 Agent 在 `MAX_STEPS`(默认为 12)内没有采取终止行动,它将被强制执行上报操作,而不是无限循环下去。
## 适配到真实系统
只需更改 `agent.py` 中这六个只读函数的**函数体**,将其指向真实的环境,即可完成适配:
| 函数 | 在真实部署中将调用 |
|---|---|
| `get_failed_task_logs` | 您的编排器 API(Airflow, Databricks Workflows 等) |
| `get_recent_schema_changes` | 您的元数据/schema 目录 |
| `run_diagnostic_query` | 您的数据仓库,通过严格只读的凭证进行访问 |
| `get_upstream_task_status` | 您的编排器 API |
| `get_recent_code_changes` | 您的版本控制系统的 API |
| `search_past_incidents` | 您的事件跟踪系统,或记录以往解决方案的简单日志 |
工具的**schema**、**系统提示词**和**Agent 循环**完全不需要更改——这正是以这种方式设计工具边界的意义所在。此外,您至少还需要做到以下几点:
- 为每个工具的后端凭证添加真实的身份验证/授权范围限制,且范围应尽可能小
- 添加对每次工具调用和结果的日志记录以备日后审计,而不仅仅是输出到控制台
- 将 `propose_fix` 和 `escalate_to_human` 路由到您的团队实际进行审查的地方(Slack、事件管理工具、电子邮件),而不是仅仅打印到控制台上
- 在将其应用于任何线上环境之前,构建一个小型包含过往事件的评估集,并针对这些事件对 Agent 进行实际测试——有关评估、模型选择和失败模式的完整讨论,请参阅 [ARTICLE.md](ARTICLE.md)。
## 许可证
MIT — 查看 [LICENSE](LICENSE)。
标签:AI智能体, AI风险缓解, C2, DLL 劫持, LLM评估, Ollama, Python, 大语言模型, 数据工程, 无后门, 自动化运维, 逆向工具