rkalra247/data-incident-commander
GitHub: rkalra247/data-incident-commander
基于 DataHub 的 AI 数据事件响应指挥官,将断言失败和 schema 变更等信号自动转化为从检测、根因排查到复盘报告的完整事件处理流程。
Stars: 0 | Forks: 0
# 数据事件指挥官
**在仪表盘出现错误之前,将无声的断言失败转化为一个已诊断、已认领、已解决的数据事件。**
面向**数据平台**的 AI 事件指挥官。当一个 DataHub 信号——失败的断言、schema 变更、中断的 pipeline 运行——到来时,无需人工值班人员开启缓慢、手动的排查(哪个数据集?谁拥有它?上游发生了什么变更?下游会受什么影响?),Agent 会接管整个事件:
1. **检测并分类**将信号转化为结构化的数据事件(严重程度、DataHub 类型、失败资产 URN、置信度)。
2. **开启数据事件作战室**并召集相关人员——失败资产的所有者和 data steward(来自 DataHub 的所有权信息)、可疑上游变更的提交者,以及一名业务关键的下游消费者。
3. **收集 lineage 上下文**——下游的**影响范围**(数据集、仪表盘、图表及其所有者),断言运行历史(失败的检查*以及*那些保持通过状态从而导致损失悄无声息发生的检查),上游变更元数据,以及同一资产上的过往事件——然后综合生成一份简报。
4. **指派所有者**并附带推理逻辑(通常是上游变更的提交者)。
5. **排列根因假设**,随着每个数据团队提供其掌握的关键事实,假设的概率会不断精确。
6. **推荐补救措施**,跟踪实时状态仪表盘,并为事件作战室**以及**下游数据消费者起草更新通知。
7. **在解决时自动撰写复盘报告**,包含影响摘要:*在 N 分钟内解决 · 节省了工程 + 分析师的小时数 · 遏制了下游风险暴露 · 上下文获取时间缩短约 78%*——并且可以将解决方案**写回 DataHub**,以便下次发生时能自动诊断。
它是**作战室 Agent** 事件响应主干(检测 → 开启作战室 → 收集上下文 → 指派所有者 → 排列根因 → 补救 → 通讯 → 复盘)在数据领域的重新应用:这里的事件是失败的数据集/断言,而影响是 **lineage 影响范围**,而不是某个客户账户。
## 真实与模拟的界限
本项目采用**离线优先**发布,因此整个演示无需任何密钥和外部设置即可运行。
| 功能 | 默认(离线)| 真实模式 |
|---|---|---|
| **LLM / agents** | `MockBrain` — 确定性,无需 API 密钥 | `AnthropicBrain` — 真实的 Claude (`claude-opus-4-8`) |
| **DataHub 元数据**(lineage、断言、所有权、历史、经验)| `MockDataHub` — 一个预设场景 (`server/integrations/mock.ts`) | `DataHubGraphQL` — 通过 `POST /api/graphql` + PAT 连接的实时 GMS |
| **事件作战室聊天** | `SimSlackProvider` — 浏览器内的作战室 UI | `RealSlackProvider` — 通过 Web API 连接的真实工作区 |
模拟 LLM 的文本撰写得就像真实的 Claude 对 lineage 影响范围的综合分析,因此离线演示本身就极具说服力。真实模式是按子系统选择启用的,并且 DataHub GraphQL provider **在任何单次调用失败时都会回退到模拟模式**,因此实时演示绝不会发生硬故障。
**诚实的边界。** DataHub 的读取 API(lineage、断言、所有权、事件)和 Actions Framework 事件流是 OSS;自动运行断言*监控*以及 Slack/webhook *订阅*是 Acryl Cloud 的一部分。因此我们**模拟了触发器**——`POST /api/datahub/event` 接受符合 DataHub 真实数据结构的合成 `AssertionRunEvent`/`EntityChangeEvent`,并驱动与实时 webhook 相同的检测逻辑。真实的 `raiseIncident` / `updateIncidentStatus` 写入路径以及 `searchAcrossLineage` 读取路径均已针对真实的 GraphQL schema 实现(`server/integrations/datahub-graphql.ts`)。
## 快速开始
```
npm run setup # install server + web deps
npm run demo # print the full offline incident flow to stdout (no server, no keys)
```
`npm run demo` 将端到端演示预设的 **`fct_transactions` 无声计数偏低**事件——检测 → 作战室 → lineage 简报 → 所有者 → 根因排序 → 补救 → 状态仪表盘 → 通讯 → 复盘——并打印出记录。这是了解产品功能的最快方式。
## 运行服务端 + 仪表盘
```
npm run dev:all # API on :4000, web dashboard on :5173 (Vite)
```
- **仅 API:** `npm run dev`(`tsx watch`,端口 `4000`,可通过 `PORT` 覆盖)。
- **仅 Web:** `npm run web:dev`。
- **类型检查:** `npm run typecheck` · **测试:** `npm test`。
打开仪表盘,点击 **Run scripted incident**,在数据事件作战室 + 事件详情(影响范围、断言时间线、排列后的根因、复盘报告)中实时观看 `fct_transactions` 流程演练——无需终端。每个视图都支持深链接(`/incidents/:id`)。
## Provider 开关
将 `.env.example` 复制到 `.env` 并设置相应的开关,即可从模拟模式切换到真实模式。**如果不设置任何内容,一切将在离线状态下运行。**
```
# 真正的 Claude 而非 mock brain
LLM_PROVIDER=anthropic
ANTHROPIC_API_KEY=sk-ant-...
ANTHROPIC_MODEL=claude-opus-4-8 # default
# 真正的 DataHub 而非 mock metadata
DATAHUB_PROVIDER=graphql
DATAHUB_GMS_URL=http://localhost:8080 # your GMS
DATAHUB_GMS_TOKEN=... # PAT: DataHub → Settings → Access Tokens
# 真正的 Slack workspace 而非模拟的 room
SLACK_PROVIDER=slack
SLACK_BOT_TOKEN=xoxb-... # scopes: channels:manage, channels:write, chat:write, groups:write, users:read
```
如需使用真实的 DataHub 写入路径(`raiseIncident`),请在**本地运行 DataHub Core**(使用 docker quickstart)并生成一个 PAT——`demo.datahub.com` 仅支持浏览。请参阅 [`docs/ARCHITECTURE.md`](docs/ARCHITECTURE.md) 了解 Agent 列表、DataHub 集成接口,以及端到端断言失败事件的请求追踪。
## 像真实的 DataHub 一样触发检测
`POST /api/datahub/event` 映射了 DataHub 的 Actions Framework webhook——发送一个合成的断言运行记录,指挥官就会启动:
```
curl -s localhost:4000/api/datahub/event -H 'content-type: application/json' -d '{
"asserteeUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,analytics.core.fct_transactions,PROD)",
"assertionInfo": { "type": "VOLUME", "description": "Row count within 1% of upstream source" },
"result": { "type": "FAILURE", "nativeResults": [{ "key": "missingRows", "value": "43120" }] }
}'
```
## 演示脚本(`fct_transactions` 流程)
预设的事件(`server/demo/scenario.ts`)是一个棘手的、跨团队的数据事件:没有 pipeline 错误,没有失败的任务——一个**无声的计数偏低**,只有 DataHub 的 VOLUME 断言能捕捉到,而 FRESHNESS 断言却保持通过(绿色)。
1. 一个 DataHub **VOLUME 断言失败**,发生在 `analytics.core.fct_transactions`(Snowflake PROD)上:行数比上游数据源低约 2%——在 60 分钟内静默丢失了约 4.3 万条交易记录,日志中没有错误。它为高管收入仪表盘和监管对账提供数据。
2. **Detection Agent** 对其进行分类:`fct_transactions` 上的 **SEV-1 volume** 事件。
3. 一个**数据事件作战室**开启;邀请相关成员(资产所有者 + data steward + 上游变更提交者 + 一名业务关键的下游消费者 + 领导层 [针对 SEV-1])。
4. **Context Agent** 发布简报:**影响范围**(下游的高管 `daily_revenue` mart 和对账报告),失败的 VOLUME 检查与通过的 FRESHNESS 检查对比,可疑的上游变更,以及匹配的历史事件 **INC-207**。
5. **Commander** 将上游变更的提交者指派为所有者。
6. **Root Cause Agent** 发布主要假设(约 63%):上游的数据接入变更静默丢弃了数据行。
7. 补救清单 + 实时**状态仪表盘**(SEV-1,缓解中,下游风险暴露)。
8. **Comms Agent** 起草作战室更新和下游消费者通知。
9. 作战室成员各司其职——分析工程团队将数据缺口范围缩小到某个区域故障转移窗口,另一人审计确认 dbt 转换无误,data steward 发现了特定区域的去重/拓扑不匹配问题,接入工程师确认他的自适应批处理变更是原因(与 INC-207 形态相同)→主要假设概率攀升至 **~84%**,并得到四个信号的佐证;状态 → 监控中。
10. **已解决。** 丢弃的事件已重播,方差恢复为零。复盘报告自动生成;解决方案被写回 DataHub,以便下次发生时自动诊断。
## 目录布局
```
server/
agents/ classify (deterministic), brain interface, mock-brain, anthropic-brain
services/ orchestrator (WarRoom), context (lineage/blast-radius), exec, app (composition root)
integrations/ DataHub adapter interface + mock, real GraphQL client, Actions-webhook event ingest
slack/ provider interface, sim (in-browser room), real (@slack/web-api)
db/ SQLite store + schema (incidents, participants, timelines, action_items, hypotheses)
api/ HTTP surface + SSE event stream
demo/ scripted fct_transactions scenario + headless runner
web/ Vite + React + Tailwind dashboard (data-incident room + incident detail)
docs/ architecture, build plan
```
请参阅 [`AGENTS.md`](AGENTS.md) 了解开发约定以及在进行代码扩展时各个模块的位置。
## 黑客松参赛说明
**灵感来源。** 数据事件响应过程缓慢且依赖手动。当一个数据集悄无声息地损坏时,能够快速解决它的上下文信息其实已经存在于 DataHub 中——lineage、所有权、断言历史、过往事件——但没有任何东西去*利用*它们。值班人员仍然需要注意到警报,弄清楚是哪个数据集,追踪是谁拥有的,排查上游发生了什么变更,理清下游会受什么影响,并把合适的人拉到一个房间里。这意味着在错误的数字显现在高管仪表盘上的同时,我们要白白耗费 30-60 分钟去收集上下文。我们之前已经构建了 **作战室 Agent**——一个针对客户升级事件的 AI 事件指挥官——我们意识到完全相同的主干逻辑同样适用于数据事件,只是这里的“影响范围”变成了 lineage 图,而不是客户账户。
**功能介绍。** 它将原始的 DataHub 信号转化为一个全面运作的事件:检测并分类 → 开启一个有专职人员的数据事件作战室 → 汇总 lineage 影响范围简报 → 根据 DataHub 所有权和上游变更元数据指派所有者 → 随着各数据团队的介入,不断精确根因假设的排序 → 推荐补救措施 → 为事件作战室和下游消费者起草通讯稿 → 自动撰写无指责的复盘报告,并将解决方案写回 DataHub,以便下一次出现相同故障时能自动诊断。整个过程会实时流式传输到仪表盘上。
**构建方式。** 一个单一的 Node + TypeScript (ESM) Express 服务管理着进程内事件总线、作为唯一事实来源的 SQLite 存储,以及位于一个与供应商无关的 `Brain` 接口(包含模拟和真实 Claude)背后的所有 AI 能力。DataHub 位于一个 `DataHub` 适配器接口之后,配有一个确定性的模拟实现(用于离线演示)和一个真实的 GraphQL 客户端(`searchAcrossLineage` 用于获取影响范围,`raiseIncident`/`updateIncidentStatus` 用于写入路径,断言 `runEvents` 和所有权用于获取上下文)。`POST /api/datahub/event` 数据接入映射了 DataHub 的 Actions Framework webhook,使得离线演示和真实的 DataHub 能够触发完全相同的 pipeline。React + Vite + Tailwind 仪表盘订阅了 Server-Sent-Events 流,无需轮询。真实与模拟的选择完全在一个 composition root 中决定;真实的 DataHub 路径在遇到任何失败时都会回退到模拟模式。
**面临的挑战。** 最难处理的事件是那些*无声*的事件——没有错误,没有失败的任务,只有一个单纯出错的数据集——因此演示场景设定为 2% 的计数偏低,只有在 FRESHNESS 保持绿色的同时 VOLUME 断言才能捕捉到它,而且只有当四个团队各自提供一个事实时才能解决。保持模拟 LLM 的确定性*并且*极具说服力(让离线演示读起来像真实的分析)需要扎实的编写功底。此外,DataHub 的功能接口是割裂的——读取和事件框架是 OSS,但自动运行监控和订阅是 Cloud 的——因此我们划定了一条清晰的界限:模拟触发器,针对真实的 GraphQL schema 实现实际的读/写路径,并确保真实模式永不发生硬故障。
**未来计划。** 将 OSS Actions Framework Kafka sink 直接连接到事件接入端;通过官方的 `mcp-server-datahub` 丰富上下文;扩展学习闭环,使写回的解决方案在整个 lineage 图中不断产生复合效应。
**技术栈:** TypeScript · Node 22 · Express · SQLite (better-sqlite3) · Server-Sent Events · React · Vite · Tailwind · Zod · Anthropic SDK (Claude `claude-opus-4-8`) · DataHub GraphQL API + Actions Framework。
## 许可证
[Apache-2.0](LICENSE)。
标签:AI智能体, DataHub, MITM代理, 数据治理, 数据血缘, 数据运维, 自动化攻击, 自动化运维