第321篇:网络知识图谱与 AI Ops

关键词

知识图谱、AI Ops、智能运维、图数据库、Neo4j、拓扑推理、根因分析、知识推理


一、网络知识图谱概述

1.1 为什么需要知识图谱

传统网络管理的知识困境:

  知识分散:
  ┌─ 设备配置在 CLI 中
  ├─ 拓扑信息在图纸中
  ├─ 故障案例在文档中
  ├─ 专家经验在工程师脑中
  └─ 告警在监控系统中

  问题:
  ┌─ 信息孤岛,难以关联分析
  ├─ 排障依赖个人经验
  ├─ 知识难以传承
  └─ 无法推理隐含关系

  知识图谱的解决思路:
  ┌─ 将网络知识建模为"实体-关系"图
  ├─ 设备、接口、VLAN、路由协议都是节点
  ├─ 连接、承载、依赖都是关系
  └─ 在图结构上推理和分析

1.2 网络知识图谱模型

网络知识图谱的实体和关系:

实体类型: ┌─ 物理实体:设备、板卡、接口、光模块 ├─ 逻辑实体:VLAN、VRF、BGP 进程、OSPF 区域 ├─ 协议实体:BGP Peer、OSPF Neighbor、ISIS Adjacency ├─ 业务实体:租户、站点、业务、应用 └─ 事件实体:告警、变更、故障、日志

关系类型: ┌─ 物理连接:设备-连接-设备(光纤/网线) ├─ 逻辑连接:接口-属于-设备、VLAN-包含-接口 ├─ 协议关系:BGP Peer(设备A-对等-设备B) ├─ 依赖关系:业务-依赖-VPN、VPN-依赖-接口 ├─ 承载关系:VLAN-承载-业务 └─ 影响关系:告警-影响-设备、故障-影响-业务

示例: | [CORE-SW01] ├──[GE0/0/1]──连接──[GE0/0/24]──属于──[ACC-SW01] └──[BGP_Peer]──对等──[CORE-RT01] | 属于 └──[OSPF_Neighbor]→[ACC-SW01] | IP: 10.0.1.1 VLAN: [10]──包含──[VLAN10] 承载 | | --- | --- | --- |


二、Neo4j 图数据库构建

2.1 安装与初始化

#!/usr/bin/env python3
# neo4j_network_graph.py — Neo4j 网络知识图谱

# 安装:python -m pip install neo4j

from neo4j import GraphDatabase
import json

class NetworkKnowledgeGraph:
    """网络知识图谱"""

    def __init__(self, uri="bolt://localhost:7687",
                 user="neo4j", password="password"):
        self.driver = GraphDatabase.driver(uri, auth=(user, password))

    def close(self):
        self.driver.close()

    def clear(self):
        """清空图谱"""
        with self.driver.session() as session:
            session.run("MATCH (n) DETACH DELETE n")
        print("✽ 图谱已清空")

    # ===== 创建设备节点 =====
    def create_device(self, name: str, role: str,
                      vendor: str, model: str, mgmt_ip: str):
        """创建设备节点"""
        with self.driver.session() as session:
            result = session.run(
                """
                MERGE (d:Device {name: $name})
                SET d.role = $role,
                    d.vendor = $vendor,
                    d.model = $model,
                    d.mgmt_ip = $mgmt_ip
                RETURN d
                """,
                name=name, role=role, vendor=vendor,
                model=model, mgmt_ip=mgmt_ip,
            )
            return result.single()

    # ===== 创建接口节点 =====
    def create_interface(self, device_name: str, if_name: str,
                         ip: str = "", mac: str = "",
                         status: str = "down"):
        """创建接口节点并关联设备"""
        with self.driver.session() as session:
            session.run(
                """
                MATCH (d:Device {name: $device_name})
                MERGE (i:Interface {name: $if_name, device: $device_name})
                SET i.ip = $ip,
                    i.mac = $mac,
                    i.status = $status
                MERGE (i)-[:BELONGS_TO]->(d)
                """,
                device_name=device_name, if_name=if_name,
                ip=ip, mac=mac, status=status,
            )

    # ===== 创建物理连接 =====
    def create_link(self, dev_a: str, intf_a: str,
                    dev_b: str, intf_b: str,
                    speed: str = "", status: str = "up"):
        """创建设备间物理连接"""
        with self.driver.session() as session:
            session.run(
                """
                MATCH (ia:Interface {name: $intf_a, device: $dev_a})
                MATCH (ib:Interface {name: $intf_b, device: $dev_b})
                MERGE (ia)-[l:CONNECTS {speed: $speed, status: $status}]->(ib)
                SET l.speed = $speed, l.status = $status
                """,
                dev_a=dev_a, intf_a=intf_a,
                dev_b=dev_b, intf_b=intf_b,
                speed=speed, status=status,
            )

    # ===== 创建 VLAN 节点 =====
    def create_vlan(self, vlan_id: int, name: str):
        """创建 VLAN 节点"""
        with self.driver.session() as session:
            session.run(
                """
                MERGE (v:VLAN {id: $vlan_id})
                SET v.name = $name
                """,
                vlan_id=vlan_id, name=name,
            )

    def assign_vlan_to_interface(self, device_name: str,
                                  if_name: str, vlan_id: int):
        """关联 VLAN 到接口"""
        with self.driver.session() as session:
            session.run(
                """
                MATCH (i:Interface {name: $if_name, device: $device_name})
                MATCH (v:VLAN {id: $vlan_id})
                MERGE (i)-[:CARRIES]->(v)
                """,
                device_name=device_name,
                if_name=if_name,
                vlan_id=vlan_id,
            )

    # ===== 创建 BGP 邻居关系 =====
    def create_bgp_peer(self, dev_a: str, local_as: int,
                        dev_b: str, remote_as: int, state: str = "established"):
        """创建 BGP 邻居关系"""
        peer_name = f"{dev_a}_BGP_{dev_b}"
        with self.driver.session() as session:
            session.run(
                """
                MATCH (da:Device {name: $dev_a})
                MATCH (db:Device {name: $dev_b})
                MERGE (da)-[p:BGP_PEER {name: $peer_name}]->(db)
                SET p.local_as = $local_as,
                    p.remote_as = $remote_as,
                    p.state = $state
                """,
                dev_a=dev_a, peer_name=peer_name,
                dev_b=dev_b, local_as=local_as,
                remote_as=remote_as, state=state,
            )

    # ===== 查询 =====
    def query_device(self, name: str) -> dict:
        """查询设备信息"""
        with self.driver.session() as session:
            result = session.run(
                "MATCH (d:Device {name: $name}) RETURN d",
                name=name,
            )
            record = result.single()
            return dict(record["d"]) if record else {}

    def query_device_connections(self, device_name: str) -> list:
        """查询设备的所有连接"""
        with self.driver.session() as session:
            result = session.run(
                """
                MATCH (d:Device {name: $device_name})
                MATCH (d)<-[:BELONGS_TO]-(i:Interface)
                MATCH (i)-[c:CONNECTS]->(other_intf:Interface)
                MATCH (other_intf)-[:BELONGS_TO]->(other_dev:Device)
                RETURN i.name AS local_intf, c.speed AS speed,
                       other_intf.name AS remote_intf,
                       other_dev.name AS remote_device
                """,
                device_name=device_name,
            )
            return [dict(r) for r in result]

    def query_path(self, start: str, end: str, max_hops: int = 10) -> list:
        """查询设备间路径"""
        with self.driver.session() as session:
            result = session.run(
                """
                MATCH path = shortestPath(
                    (d1:Device {name: $start})
                    -[:BELONGS_TO|CONNECTS*..$max_hops]-
                    (d2:Device {name: $end})
                )
                RETURN [n in nodes(path) | n.name] AS device_path
                """,
                start=start, end=end, max_hops=max_hops,
            )
            return [dict(r) for r in result]

    def query_impact_analysis(self, device_name: str) -> list:
        """影响面分析:设备故障会影响到什么"""
        with self.driver.session() as session:
            result = session.run(
                """
                MATCH (d:Device {name: $device_name})
                OPTIONAL MATCH (d)-[:BGP_PEER]->(bgp_dev:Device)
                OPTIONAL MATCH (d)<-[:BELONGS_TO]-(:Interface)
                    -[:CARRIES]->(v:VLAN)
                RETURN d.name AS device,
                       collect(DISTINCT bgp_dev.name) AS bgp_affected,
                       collect(DISTINCT v.name) AS vlans_affected
                """,
                device_name=device_name,
            )
            return [dict(r) for r in result]


# ===== 构建图谱 =====
def build_example_graph():
    """构建示例网络图谱"""
    kg = NetworkKnowledgeGraph()
    kg.clear()

    # 创建设备
    devices = [
        ("CORE-SW01", "core", "Huawei", "CE12808", "10.0.0.1"),
        ("CORE-RT01", "core", "Huawei", "NE20E-S4", "10.0.0.2"),
        ("ACC-SW01", "access", "Huawei", "S5735", "10.0.1.1"),
        ("ACC-SW02", "access", "Huawei", "S5735", "10.0.1.2"),
    ]
    for name, role, vendor, model, ip in devices:
        kg.create_device(name, role, vendor, model, ip)
    print("✓ 设备节点创建完成")

    # 创建接口
    kg.create_interface("CORE-SW01", "40GE0/0/1", "10.0.12.1")
    kg.create_interface("CORE-RT01", "GE0/0/1", "10.0.12.2")
    kg.create_interface("CORE-SW01", "GE0/0/1", "10.0.13.1")
    kg.create_interface("ACC-SW01", "GE0/0/1", "10.0.13.2")
    kg.create_interface("CORE-SW01", "GE0/0/2")
    kg.create_interface("ACC-SW02", "GE0/0/1")
    kg.create_interface("ACC-SW01", "GE0/0/2", "10.0.20.1")
    kg.create_interface("ACC-SW02", "GE0/0/2", "10.0.20.2")
    print("✓ 接口节点创建完成")

    # 创建连接
    kg.create_link("CORE-SW01", "40GE0/0/1", "CORE-RT01", "GE0/0/1", "40G")
    kg.create_link("CORE-SW01", "GE0/0/1", "ACC-SW01", "GE0/0/1", "10G")
    kg.create_link("CORE-SW01", "GE0/0/2", "ACC-SW02", "GE0/0/1", "10G")
    kg.create_link("ACC-SW01", "GE0/0/2", "ACC-SW02", "GE0/0/2", "10G")
    print("✓ 物理连接创建完成")

    # 创建 VLAN
    for vid, name in [(10, "MGMT"), (20, "OFFICE"), (30, "GUEST")]:
        kg.create_vlan(vid, name)
    print("✓ VLAN 节点创建完成")

    # 关联 VLAN 到接口
    kg.assign_vlan_to_interface("ACC-SW01", "GE0/0/1", 10)
    kg.assign_vlan_to_interface("ACC-SW01", "GE0/0/1", 20)
    kg.assign_vlan_to_interface("ACC-SW02", "GE0/0/1", 10)
    print("✓ VLAN 关联完成")

    # 创建 BGP 邻居
    kg.create_bgp_peer("CORE-SW01", 65001, "CORE-RT01", 65002, "established")
    print("✓ BGP 关系创建完成")

    # 查询示例
    print(f"\n=== 查询 CORE-SW01 ===")
    info = kg.query_device("CORE-SW01")
    print(f"  角色: {info.get('role')}, 型号: {info.get('model')}")

    print(f"\n=== CORE-SW01 的连接 ===")
    for conn in kg.query_device_connections("CORE-SW01"):
        print(f"  {conn['local_intf']} ({conn['speed']}) → "
              f"{conn['remote_device']} {conn['remote_intf']}")

    print(f"\n=== 路径查询: CORE-RT01 → ACC-SW02 ===")
    for path in kg.query_path("CORE-RT01", "ACC-SW02"):
        print(f"  路径: {' → '.join(path['device_path'])}")

    print(f"\n=== CORE-SW01 故障影响面 ===")
    for impact in kg.query_impact_analysis("CORE-SW01"):
        print(f"  BGP 影响: {impact['bgp_affected']}")
        print(f"  VLAN 影响: {impact['vlans_affected']}")

    kg.close()
    return kg

if __name__ == "__main__":
    build_example_graph()

三、AI Ops 核心能力

3.1 异常检测

#!/usr/bin/env python3
# aiops_anomaly.py — AI Ops 异常检测

"""
基于统计和 ML 的异常检测
"""

import numpy as np
from collections import deque
from datetime import datetime, timedelta

class AnomalyDetector:
    """异常检测器"""

    def __init__(self, window_size: int = 10, threshold: float = 2.0):
        self.window = deque(maxlen=window_size)
        self.threshold = threshold

    def add_point(self, value: float) -> bool:
        """添加数据点,返回是否异常"""
        self.window.append(value)

        if len(self.window) < 5:  # 数据不足
            return False

        # Z-Score 检测
        mean = np.mean(self.window)
        std = np.std(self.window)

        if std == 0:
            return False

        z_score = abs(value - mean) / std
        is_anomaly = z_score > self.threshold

        if is_anomaly:
            print(f"⚠ 异常检测: {value:.2f}, Z-Score={z_score:.2f}, "
                  f"均值={mean:.2f}, 标准差={std:.2f}")

        return is_anomaly


class BaselineAnomalyDetector:
    """基于基线的异常检测(适用于周期性指标)"""

    def __init__(self):
        self.baselines = {}  # metric → {hour → [values]}

    def learn(self, metric: str, value: float, timestamp: datetime):
        """学习基线"""
        hour = timestamp.hour
        if metric not in self.baselines:
            self.baselines[metric] = {}
        if hour not in self.baselines[metric]:
            self.baselines[metric][hour] = []

        self.baselines[metric][hour].append(value)

    def detect(self, metric: str, value: float,
               timestamp: datetime) -> bool:
        """检测异常"""
        hour = timestamp.hour
        if metric not in self.baselines or hour not in self.baselines[metric]:
            return False

        values = self.baselines[metric][hour]
        if len(values) < 5:
            return False

        mean = np.mean(values)
        std = np.std(values)

        if std == 0:
            return value != mean

        z_score = abs(value - mean) / std
        return z_score > 3.0  # 3σ

# 测试
detector = BaselineAnomalyDetector()
# 学习一周的 CPU 数据
for day in range(7):
    for hour in range(24):
        # 正常工作负载
        cpu = 30 + 10 * np.sin(hour / 24 * 2 * np.pi)
        detector.learn(
            "cpu_usage", cpu,
            datetime(2025, 1, 1) + timedelta(days=day, hours=hour),
        )

# 检测新的数据点
now = datetime(2025, 1, 8, 14, 0)
test_values = [35, 38, 40, 85, 92, 45, 42]
for v in test_values:
    if detector.detect("cpu_usage", v, now):
        print(f"⚠ 异常CPU: {v}% (时间: {now.hour}:00)")
    else:
        print(f"✓ 正常CPU: {v}%")

3.2 关联分析

#!/usr/bin/env python3
# aiops_correlation.py — 告警关联分析

from collections import defaultdict
from datetime import datetime, timedelta
import json

class AlarmCorrelation:
    """告警关联分析"""

    def __init__(self, time_window: int = 300):
        self.time_window = time_window  # 关联时间窗口(秒)
        self.alarms = []

    def add_alarm(self, device: str, alarm_type: str,
                  severity: str, timestamp: datetime,
                  message: str = ""):
        """添加告警"""
        self.alarms.append({
            "device": device,
            "type": alarm_type,
            "severity": severity,
            "timestamp": timestamp,
            "message": message,
        })

    def correlate(self) -> list:
        """执行关联分析"""
        # 按时间排序
        sorted_alarms = sorted(self.alarms, key=lambda a: a["timestamp"])

        groups = []
        current_group = []

        for alarm in sorted_alarms:
            if not current_group:
                current_group.append(alarm)
                continue

            # 检查时间窗口
            time_diff = (
                alarm["timestamp"] - current_group[0]["timestamp"]
            ).total_seconds()

            if time_diff <= self.time_window:
                current_group.append(alarm)
            else:
                if len(current_group) >= 2:
                    groups.append(self._analyze_group(current_group))
                current_group = [alarm]

        if len(current_group) >= 2:
            groups.append(self._analyze_group(current_group))

        return groups

    def _analyze_group(self, group: list) -> dict:
        """分析告警组,识别根因"""
        devices = set(a["device"] for a in group)
        types = set(a["type"] for a in group)

        # 根因推断
        root_cause = None
        dependencies = []

        if "LINK_DOWN" in types:
            # 链路 DOWN 可能是根因
            root_cause = next(
                a for a in group if a["type"] == "LINK_DOWN"
            )
        elif "BGP_DOWN" in types:
            root_cause = next(
                a for a in group if a["type"] == "BGP_DOWN"
            )
        elif "DEVICE_DOWN" in types:
            root_cause = next(
                a for a in group if a["type"] == "DEVICE_DOWN"
            )

        return {
            "time_range": f"{group[0]['timestamp']} ~ {group[-1]['timestamp']}",
            "alarm_count": len(group),
            "devices": list(devices),
            "root_cause": {
                "device": root_cause["device"] if root_cause else "unknown",
                "type": root_cause["type"] if root_cause else "unknown",
                "message": root_cause["message"] if root_cause else "",
            } if root_cause else None,
            "alarms": group,
        }

# 测试
corr = AlarmCorrelation(time_window=300)

# 模拟告警序列:链路 DOWN → BGP DOWN → 路由变化
now = datetime.now()
corr.add_alarm("CORE-SW01", "LINK_DOWN", "critical",
               now, "40GE0/0/1 link down")
corr.add_alarm("CORE-RT01", "BGP_DOWN", "major",
               now + timedelta(seconds=10), "BGP peer 10.0.12.1 down")
corr.add_alarm("CORE-SW01", "ROUTE_CHANGE", "warning",
               now + timedelta(seconds=30), "Route to 10.0.0.0/16 withdrawn")
corr.add_alarm("ACC-SW01", "PING_LOSS", "minor",
               now + timedelta(seconds=60), "Ping loss to 10.0.0.1")

groups = corr.correlate()
for g in groups:
    print(f"=== 告警组 ({g['alarm_count']} 条) ===")
    print(f"时间: {g['time_range']}")
    print(f"设备: {', '.join(g['devices'])}")
    if g['root_cause']:
        print(f"根因: {g['root_cause']['device']} - {g['root_cause']['type']}")
    else:
        print("根因: 未识别")
    print()

四、AI Ops 成熟度演进

AI Ops 四层演进模型:

  L1: 被动监控
  ┌─ 告警阈值判断(CPU > 80% → 告警)
  ├─ 人工排障
  ├─ 日报周报
  └─ 基础仪表盘

  L2: 辅助分析
  ┌─ 告警压缩(风暴抑制)
  ├─ 告警关联(时间窗口关联)
  ├─ 智能基线(动态阈值)
  └─ 自动巡检报告

  L3: 智能决策
  ┌─ 根因分析(拓扑+告警关联)
  ├─ 影响评估(故障影响范围)
  ├─ 自动修复(已知故障自动处理)
  └─ 容量预测(趋势分析)

  L4: 自治网络
  ┌─ 自动规划与优化
  ├─ 预测性维护
  ├─ 零接触运维
  └─ 闭环自愈

  当前行业普遍处于 L1~L2
  领先企业进入 L3
  L4 仍在探索

五、总结

知识图谱 + AI Ops 的未来:

  知识图谱 — 网络知识的"大脑"
  ┌─ 实体关系建模:设备、链路、业务、事件
  ├─ 图推理:路径分析、影响面分析
  ├─ 根因定位:告警拓扑追溯
  └─ 知识积累:每次排障沉淀为新知识

  AI Ops — 运维智能化的"引擎"
  ┌─ 异常检测:从静态阈值到智能基线
  ├─ 关联分析:从单点告警到根因识别
  ├─ 预测分析:从被动响应到主动预防
  └─ 自动修复:从人工处理到闭环自愈

  终极目标:
  ┌─ 监控发现异常 → 图谱关联分析 → AI 定位根因
  ├─ 自动生成修复方案 → 审批执行 → 验证闭环
  └─ 知识自动沉淀 → 图谱持续更新 → 能力持续提升

下篇预告:第322篇 — 中小企业网络架构设计,将进入第8章综合案例与专家实践,通过实际案例讲解网络规划与设计。