知识图谱完整指南:概念、构建流程与 Python 实现

适用范围: 工业知识图谱(制造、设备、供应链、质量等)
技术栈: Python + Neo4j / NetworkX / RDFLib + NLP 工具


一、知识图谱简介

1.1 什么是知识图谱

知识图谱(Knowledge Graph, KG)是一种用图结构(节点-边-属性)表示客观世界实体与关系的语义网络。它由结构化的事实组成,支持语义查询与推理。

核心要素:

  • 实体(Entity/Node): 具体的对象或概念,如 “压缩机”、“设备A”、“供应商X”
  • 关系(Relation/Edge): 实体之间的语义连接,如 “隶属于”、“故障导致”、“位于”
  • 属性(Property): 实体或关系的键值对,如 “温度=85℃”、“压力=0.6MPa”

1.2 知识图谱的组成

层级 内容 说明
数据层 三元组 (头实体, 关系, 尾实体) 如 (压缩机B, 故障类型, 过热)
模式层/本体 本体(Ontology)、类(Class)、约束 定义实体类型、关系类型、继承、基数等
推理层 规则引擎、推理机 基于本体和规则进行隐式知识推导
应用层 问答、推荐、可视化、搜索 面向最终用户的接口

1.3 工业知识图谱应用场景

  • 设备维修知识库: 故障现象 → 根因 → 维修方案
  • 生产工艺优化: 工序 → 参数 → 质量结果 关联分析
  • 供应链风险洞察: 供应商 → 零部件 → 成品 依赖网络
  • 质量追溯: 批次 → 原料 → 设备 → 操作员 全链路关联
  • 数字孪生: 物理资产与其虚拟模型、传感器数据的映射

二、构建流程 Overview

┌─────────────────────────────────────────────────────────────┐
│                     知识图谱构建流程                          │
├─────────────┬─────────────┬──────────────┬──────────────────┤
│   1. 本体设计  │   2. 数据获取  │   3. 知识抽取  │   4. 知识融合   │
│  - 定义实体类型│  - 结构化      │  - 命名实体识别│  - 实体对齐     │
│  - 关系类型   │  - 半结构化    │  - 关系抽取    │  - 消歧         │
│  - 约束规则   │  - 非结构化文本│  - 属性抽取    │  - 冲突解决     │
├─────────────┼─────────────┼──────────────┼──────────────────┤
│   5. 知识表示  │   6. 存储与查询  │   7. 推理与质量  │   8. 应用集成   │
│  - RDF/Triple │  - 图数据库      │  - 一致性检查   │  - API 服务     │
│  - 图结构     │  - 图文件        │  - 规则推理     │  - 可视化       │
└─────────────┴─────────────┴──────────────┴──────────────────┘

三、详细步骤与 Python 实现

3.1 第一步:本体设计(Schema Design)

目的: 定义图谱的“骨架”——有哪些实体类型、关系类型、属性及其约束。

工业设备维护示例:

  • 实体类型(Classes):

    • Device(设备)
    • Component(部件)
    • Fault(故障)
    • Solution(维修方案)
    • Technician(维修人员)
    • SparePart(备件)
  • 关系类型(Relations):

    • HAS_PART(设备拥有部件)
    • MAY_HAVE_FAULT(部件可能故障)
    • REQUIRES_SOLUTION(故障对应方案)
    • NEEDS_SPARE(方案需要备件)
    • ASSIGNED_TO(任务指派人员)

Python 描述(使用 simple-ontology 类):

# ontology.py
from dataclasses import dataclass, field
from typing import List, Dict, Optional

@dataclass
class Property:
    name: str
    dtype: str  # 'string', 'integer', 'float', 'datetime'
    description: str = ""
    required: bool = False

@dataclass
class EntityClass:
    name: str
    label: str  # 显示名称
    properties: List[Property] = field(default_factory=list)
    super_classes: List[str] = field(default_factory=list)  # 继承

@dataclass
class Relation:
    name: str
    domain: str  # 头实体类型
    range: str   # 尾实体类型
    properties: List[Property] = field(default_factory=list)
    symmetric: bool = False
    transitive: bool = False

class Ontology:
    def __init__(self):
        self.classes: Dict[str, EntityClass] = {}
        self.relations: Dict[str, Relation] = {}

    def add_class(self, cls: EntityClass):
        self.classes[cls.name] = cls

    def add_relation(self, rel: Relation):
        self.relations[rel.name] = rel

# 定义本体
onto = Ontology()

# 实体类
onto.add_class(EntityClass("Device", "设备", [
    Property("name", "string", required=True),
    Property("model", "string"),
    Property("manufacturer", "string"),
    Property("install_date", "datetime")
]))

onto.add_class(EntityClass("Component", "部件", [
    Property("name", "string", required=True),
    Property("code", "string"),
    Property("spec", "string")
]))

onto.add_class(EntityClass("Fault", "故障", [
    Property("name", "string", required=True),
    Property("severity", "string"),  # 高/中/低
    Property("symptom", "string")
]))

onto.add_class(EntityClass("Solution", "维修方案", [
    Property("name", "string", required=True),
    Property("steps", "string"),  # JSON 或文本
    Property("duration", "integer", "预计分钟数")
]))

onto.add_class(EntityClass("Technician", "维修人员", [
    Property("name", "string", required=True),
    Property("level", "string"),  # 初级/中级/高级
    Property("phone", "string")
]))

onto.add_class(EntityClass("SparePart", "备件", [
    Property("name", "string", required=True),
    Property("part_number", "string"),
    Property("stock", "integer")
]))

# 关系
onto.add_relation(Relation("HAS_PART", "Device", "Component"))
onto.add_relation(Relation("MAY_HAVE_FAULT", "Component", "Fault"))
onto.add_relation(Relation("REQUIRES_SOLUTION", "Fault", "Solution"))
onto.add_relation(Relation("NEEDS_SPARE", "Solution", "SparePart"))
onto.add_relation(Relation("ASSIGNED_TO", "Solution", "Technician"))

print(f"Ontology loaded: {len(onto.classes)} classes, {len(onto.relations)} relations")

3.2 第二步:数据获取与清洗

数据源分类:

类型 示例 获取方式
结构化 MySQL 设备台账、MES 工单 SQL 查询、API
半结构化 XML 设备日志、JSON 配置 解析转换
非结构化 Word 维修手册、PDF 报告、邮件 NLP 抽取

示例:从 MySQL 拉取设备信息

# data_loader.py
import pandas as pd
import mysql.connector
from typing import List, Dict

def load_devices(mysql_config: Dict) -> List[Dict]:
    """从设备管理表加载设备实体"""
    conn = mysql.connector.connect(**mysql_config)
    query = """
        SELECT device_id, device_name, model, manufacturer, install_date
        FROM equipment.devices
        WHERE status = 'active'
    """
    df = pd.read_sql(query, conn)
    conn.close()
    return df.to_dict(orient='records')

def load_components(mysql_config: Dict) -> List[Dict]:
    """加载部件层级关系"""
    conn = mysql.connector.connect(**mysql_config)
    query = """
        SELECT component_id, component_name, code, spec, device_id
        FROM equipment.components
    """
    df = pd.read_sql(query, conn)
    conn.close()
    return df.to_dict(orient='records')

def load_faults(mysql_config: Dict) -> List[Dict]:
    """从故障报告加载故障记录"""
    conn = mysql.connector.connect(**mysql_config)
    query = """
        SELECT fault_id, component_id, fault_name, severity, symptom, report_time
        FROM maintenance.faults
    """
    df = pd.read_sql(query, conn)
    conn.close()
    return df.to_dict(orient='records')

def load_solutions(mysql_config: Dict) -> List[Dict]:
    """加载维修方案"""
    conn = mysql.connector.connect(**mysql_config)
    query = """
        SELECT solution_id, fault_id, solution_name, steps, duration, technician_id
        FROM maintenance.solutions
    """
    df = pd.read_sql(query, conn)
    conn.close()
    return df.to_dict(orient='records')

# 示例调用
mysql_cfg = {
    'host': 'localhost',
    'user': 'root',
    'password': 'your_password',
    'database': 'equipment'
}

devices = load_devices(mysql_cfg)
components = load_components(mysql_cfg)
faults = load_faults(mysql_cfg)
solutions = load_solutions(mysql_cfg)

print(f"Loaded: {len(devices)} devices, {len(components)} components, {len(faults)} faults, {len(solutions)} solutions")

3.3 第三步:知识抽取(结构化数据转三元组)

对于结构化数据,直接映射为三元组:

# triplifier.py
def device_to_triples(devices: List[Dict], components: List[Dict]) -> List[tuple]:
    """设备与部件建图"""
    triples = []
    # 设备节点
    for dev in devices:
        triples.append((
            f"device:{dev['device_id']}",
            "name",
            dev['device_name']
        ))
        triples.append((
            f"device:{dev['device_id']}",
            "model",
            dev.get('model', '')
        ))
    # 部件节点 + 关系
    for comp in components:
        triples.append((
            f"component:{comp['component_id']}",
            "name",
            comp['component_name']
        ))
        # 关系: Device -HAS_PART-> Component
        triples.append((
            f"device:{comp['device_id']}",
            "HAS_PART",
            f"component:{comp['component_id']}"
        ))
    return triples

def fault_solution_triples(faults: List[Dict], solutions: List[Dict]) -> List[tuple]:
    triples = []
    # Fault 节点
    for f in faults:
        triples.append((
            f"fault:{f['fault_id']}",
            "name",
            f['fault_name']
        ))
        triples.append((
            f"fault:{f['fault_id']}",
            "severity",
            f['severity']
        ))
        # 关系: Component -MAY_HAVE_FAULT-> Fault
        triples.append((
            f"component:{f['component_id']}",
            "MAY_HAVE_FAULT",
            f"fault:{f['fault_id']}"
        ))
    # Solution 节点 + 关系
    for sol in solutions:
        triples.append((
            f"solution:{sol['solution_id']}",
            "name",
            sol['solution_name']
        ))
        triples.append((
            f"fault:{sol['fault_id']}",
            "REQUIRES_SOLUTION",
            f"solution:{sol['solution_id']}"
        ))
    return triples

# 聚合所有三元组
device_triples = device_to_triples(devices, components)
fault_triples = fault_solution_triples(faults, solutions)
all_triples = device_triples + fault_triples
print(f"Generated {len(all_triples)} triples")

3.4 第四步:从非结构化文本抽取

如果有一些维修手册或报告在 Word/PDF 中,需要用 NLP 抽取。

方案: 使用 spaCy + 规则/机器学习进行实体识别和关系抽取。

# nlp_extractor.py
import spacy
from spacy.matcher import Matcher

# 加载中文模型(需先下载:python -m spacy download zh_core_web_sm)
nlp = spacy.load("zh_core_web_sm")

def extract_entities(text: str) -> List[tuple]:
    """抽取实体(简化版:基于词性+规则)"""
    doc = nlp(text)
    entities = []
    for ent in doc.ents:
        entities.append((ent.text, ent.label_))
    return entities

def rule_based_relation_extract(text: str):
    """规则抽取:故障-方案对"""
    doc = nlp(text)
    matcher = Matcher(nlp.vocab)
    # 模式: "X故障可通过Y方案解决"
    pattern = [
        {"TEXT": {"REGEX": ".+"}, "OP": "+"},
        {"TEXT": "故障"},
        {"TEXT": "可"},
        {"TEXT": "通过"},
        {"TEXT": {"REGEX": ".+"}, "OP": "+"},
        {"TEXT": "方案"}
    ]
    matcher.add("FAULT_SOLUTION", [pattern])
    matches = matcher(doc)
    for match_id, start, end in matches:
        span = doc[start:end]
        print("Found:", span.text)

3.5 第五步:知识融合与消歧

实体对齐(Entity Resolution): 同一实体在不同来源可能有不同名称(如 “压缩机” vs “空压机”)。

简单方法: 使用字符串相似度(Levenshtein)或余弦相似度对实体名进行聚类。

# entity_alignment.py
from thefuzz import fuzz
from itertools import combinations

def align_entities(entity_names: List[str], threshold: int = 85) -> List[List[str]]:
    """基于相似度对齐实体名,返回同一实体的别名列表"""
    clusters = []
    used = set()
    for a, b in combinations(entity_names, 2):
        if a in used or b in used:
            continue
        score = fuzz.ratio(a, b)
        if score >= threshold:
            # 找到现有簇
            found = False
            for cluster in clusters:
                if a in cluster or b in cluster:
                    cluster.add(a)
                    cluster.add(b)
                    found = True
                    break
            if not found:
                clusters.append({a, b})
            used.add(a)
            used.add(b)
    # 未匹配的单独成簇
    for name in entity_names:
        if name not in used:
            clusters.append({name})
    return [list(c) for c in clusters]

# 示例
names = ["压缩机A", "压缩机A", "压缩机B", "空压机B", "Motor", "电机"]
clusters = align_entities(names)
print("Entity clusters:", clusters)
# 输出:[['压缩机A', '压缩机A'], ['压缩机B', '空压机B'], ['Motor', '电机']]

3.6 第六步:存储与查询

选项 A:使用图数据库 Neo4j

安装 Neo4j Desktop 或 Server,使用 neo4j Python 驱动。

# neo4j_loader.py
from neo4j import GraphDatabase
import logging

class Neo4jLoader:
    def __init__(self, uri: str, user: str, password: str):
        self.driver = GraphDatabase.driver(uri, auth=(user, password))
        logging.basicConfig(level=logging.INFO)

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

    def clear(self):
        """清空数据库"""
        with self.driver.session() as session:
            session.run("MATCH (n) DETACH DELETE n")

    def create_constraints(self):
        """创建唯一约束"""
        constraints = [
            "CREATE CONSTRAINT device_id IF NOT EXISTS FOR (d:Device) REQUIRE d.id IS UNIQUE",
            "CREATE CONSTRAINT component_id IF NOT EXISTS FOR (c:Component) REQUIRE c.id IS UNIQUE",
            "CREATE CONSTRAINT fault_id IF NOT EXISTS FOR (f:Fault) REQUIRE f.id IS UNIQUE",
            "CREATE CONSTRAINT solution_id IF NOT EXISTS FOR (s:Solution) REQUIRE s.id IS UNIQUE"
        ]
        with self.driver.session() as session:
            for c in constraints:
                try:
                    session.run(c)
                    logging.info(f"Created constraint: {c}")
                except Exception as e:
                    logging.warning(f"Constraint may exist: {e}")

    def load_triples(self, triples: List[tuple]):
        """批量插入三元组(head, relation, tail)"""
        query = """
        UNWIND $batch AS row
        MERGE (h {id: row.head})
          SET h.name = row.head_name
        MERGE (t {id: row.tail})
          SET t.name = row.tail_name
        MERGE (h)-[r:`${row.rel}`]->(t)
          SET r += row.props
        """
        # 需要将三元组转换为 batch 格式
        batch = []
        for head, rel, tail in triples:
            # 解析 head/tail 的 id 和 name(根据你的命名规则)
            batch.append({
                "head": head,
                "head_name": head.split(":", 1)[1],  # 简化
                "rel": rel,
                "tail": tail,
                "tail_name": tail.split(":", 1)[1],
                "props": {}
            })
        with self.driver.session() as session:
            session.run(query, batch=batch)
            logging.info(f"Loaded {len(batch)} triples")

# 使用示例
loader = Neo4jLoader("bolt://localhost:7687", "neo4j", "password")
loader.clear()
loader.create_constraints()
loader.load_triples(all_triples[:100])  # 传你之前生成的三元组列表
loader.close()

选项 B:使用 NetworkX(本地分析/原型)

# networkx_graph.py
import networkx as nx
import matplotlib.pyplot as plt

def build_kg_nx(triples: List[tuple]) -> nx.MultiDiGraph:
    G = nx.MultiDiGraph()
    for head, rel, tail in triples:
        G.add_node(head, label=head.split(":")[0])
        G.add_node(tail, label=tail.split(":")[0])
        G.add_edge(head, tail, key=rel, relation=rel)
    return G

G = build_kg_nx(all_triples)
print(f"Graph: {G.number_of_nodes()} nodes, {G.number_of_edges()} edges")

# 简单可视化(小图)
plt.figure(figsize=(12, 8))
pos = nx.spring_layout(G, k=0.8)
nx.draw_networkx_nodes(G, pos, node_color='lightblue', node_size=500)
nx.draw_networkx_labels(G, pos, font_size=8)
edge_labels = {(u, v, k): d['relation'] for u, v, k, d in G.edges(keys=True, data=True)}
nx.draw_networkx_edge_labels(G, pos, edge_labels=edge_labels, font_size=6)
plt.axis('off')
plt.title("Industrial Maintenance Knowledge Graph")
plt.show()

3.7 第七步:知识查询与推理

Cypher 查询示例(Neo4j):

-- 1) 查找“压缩机B”可能的所有故障
MATCH (c:Component {name: '压缩机B'})-[:MAY_HAVE_FAULT]->(f:Fault)
RETURN f.name, f.severity

-- 2) 查找故障“过热”的维修方案及所需备件
MATCH (f:Fault {name: '过热'})-[:REQUIRES_SOLUTION]->(s:Solution)
MATCH (s)-[:NEEDS_SPARE]->(p:SparePart)
RETURN s.name, p.name, p.stock

-- 3) 查找能维修“过热”故障的高级技师
MATCH (f:Fault {name: '过热'})<-[:REQUIRES_SOLUTION]-(s:Solution)
MATCH (s)-[:ASSIGNED_TO]->(t:Technician)
WHERE t.level = '高级'
RETURN t.name, t.phone

-- 4) 递归查询:查找某设备的完整部件树(任意深度)
MATCH (d:Device {name: '空调机组A'})-[:HAS_PART*1..3]->(c:Component)
RETURN d, c

Python 层查询封装:

# query_engine.py
class KnowledgeQuery:
    def __init__(self, loader: Neo4jLoader):
        self.loader = loader

    def find_faults_by_component(self, component_name: str):
        query = """
        MATCH (c:Component {name: $comp})-[:MAY_HAVE_FAULT]->(f:Fault)
        RETURN f.name AS fault, f.severity AS severity
        """
        with self.loader.driver.session() as s:
            result = s.run(query, comp=component_name)
            return result.data()

    def find_solution_with_parts(self, fault_name: str):
        query = """
        MATCH (f:Fault {name: $fault})-[:REQUIRES_SOLUTION]->(s:Solution)
        MATCH (s)-[:NEEDS_SPARE]->(p:SparePart)
        RETURN s.name AS solution, collect(p.name) AS parts
        """
        with self.loader.driver.session() as s:
            return s.run(query, fault=fault_name).data()

    def find_technicians_for_fault(self, fault_name: str, level: str = None):
        lvl_cond = "WHERE t.level = $level" if level else ""
        query = f"""
        MATCH (f:Fault {{name: $fault}})<-[:REQUIRES_SOLUTION]-(s:Solution)
        MATCH (s)-[:ASSIGNED_TO]->(t:Technician)
        {lvl_cond}
        RETURN t.name AS name, t.phone AS phone
        """
        with self.loader.driver.session() as s:
            params = {"fault": fault_name}
            if level:
                params["level"] = level
            return s.run(query, **params).data()

# 使用
query = KnowledgeQuery(loader)
print(query.find_faults_by_component("压缩机A"))

3.8 第八步:推理规则

使用 Neo4j 的 APOC 或自定义规则引擎进行推理。

示例:如果备件库存低于阈值,则标记为“需采购”

CALL {
  MATCH (p:SparePart)
  WHERE p.stock < 5
  SET p.status = 'LOW_STOCK'
  RETURN p
}

复杂规则: 如果某故障频繁发生(>3次/月),则自动触发设备检查任务。

可以使用 Python 定期运行:

# rule_engine.py
from datetime import datetime, timedelta

def check_frequent_faults(threshold: int = 3):
    """最近30天内同一故障发生超过阈值,则告警"""
    query = """
    MATCH (f:Fault)<-[:REQUIRES_SOLUTION]-(s:Solution)
    WHERE s.create_time >= datetime() - duration({days: 30})
    WITH f, count(s) AS cnt
    WHERE cnt > $threshold
    RETURN f.name AS fault, cnt
    """
    with loader.driver.session() as s:
        rows = s.run(query, threshold=threshold).data()
        if rows:
            print(f"[{datetime.now()}] 高频故障告警:", rows)
            # 这里可以发送通知

并通过计划任务定时执行(如 cron / APScheduler)。

3.9 第九步:应用集成

提供 REST API:

# api_server.py
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
from typing import List, Optional

app = FastAPI(title="Industrial KG API")

class FaultSearchRequest(BaseModel):
    component: str

class SolutionResponse(BaseModel):
    fault: str
    solutions: List[str]
    parts: List[str]

@app.post("/search/faults")
def search_faults(req: FaultSearchRequest):
    faults = query.find_faults_by_component(req.component)
    return {"component": req.component, "faults": faults}

@app.get("/solution/{fault_name}")
def get_solution(fault_name: str):
    sol = query.find_solution_with_parts(fault_name)
    if not sol:
        raise HTTPException(status_code=404, detail="Fault not found")
    return sol

# 运行: uvicorn api_server:app --reload

五、评估与质量保证

指标 说明 评估方法
覆盖率 实体/关系覆盖真实世界的比例 专家抽样检查
准确率 三元组事实正确的比例 人工标注验证集
一致性 无逻辑冲突(如循环继承) 推理机检查
时效性 数据更新延迟 对比源系统时间戳

六、项目目录结构示例

industrial_kg/
├── ontology.py           # 本体定义
├── data_loader.py        # 数据库连接与原始数据加载
├── triplifier.py         # 三元组转换
├── nlp_extractor.py      # NLP 抽取(可选)
├── entity_alignment.py   # 实体对齐消歧
├── neo4j_loader.py       # Neo4j 加载器
├── query_engine.py       # 查询封装
├── api_server.py         # REST API
├── requirements.txt      # 依赖
├── config.yaml           # 配置文件(数据库连接等)
├── data/
│   ├── raw/              # 原始数据(CSV、JSON)
│   ├── processed/        # 清洗后数据
│   └── triples.csv       # 归一化三元组
├── notebooks/
│   └── exploration.ipynb # 探索性分析
└── README.md

requirements.txt:

pandas
mysql-connector-python
neo4j
networkx
matplotlib
spacy
thefuzz
fastapi
uvicorn
requests
beautifulsoup4

七、进阶话题

  • 时序知识图谱: 为三元组增加时间戳,支持时变分析
  • 向量化检索: 使用 text2vec 或 TransE 将实体嵌入,实现语义搜索
  • 本体推理: 使用 OWL 推理机(如 Apache Jena)进行自动分类
  • 增量更新: 设计 Kafka/CDC 管道监听数据库变更,实时更新图谱
  • 多模态图谱: 将图片、运维记录也作为节点嵌入

八、总结

构建一个工业知识图谱的完整流程:

  1. 明确范围 → 选定领域、实体类型
  2. 设计本体 → 定义类和关系
  3. 数据接入 → 结构化直接转三元组,非结构化用 NLP 抽取
  4. 清洗融合 → 对齐、去重、消歧
  5. 存储查询 → Neo4j(大图)或 NetworkX(原型)
  6. 质量评估 → 专家验证、准确率指标
  7. 应用开发 → 问答、搜索、可视化、告警

附录:快速运行示例

# quickstart.py
# 最小可运行:创建一个微型设备图谱并查询
from neo4j import GraphDatabase

# 1. 连接
driver = GraphDatabase.driver("bolt://localhost:7687", auth=("neo4j", "password"))

# 2. 清空
with driver.session() as s:
    s.run("MATCH (n) DETACH DELETE n")

# 3. 插入示例
triples = [
    ("device:001", "HAS_PART", "component:compA"),
    ("component:compA", "MAY_HAVE_FAULT", "fault:f1"),
    ("fault:f1", "REQUIRES_SOLUTION", "solution:s1"),
    ("solution:s1", "NEEDS_SPARE", "part:p1"),
    ("solution:s1", "ASSIGNED_TO", "tech:t1")
]

with driver.session() as s:
    for h, r, t in triples:
        s.run("""
        MERGE (a {id: $h}) SET a.name = $h_name
        MERGE (b {id: $t}) SET b.name = $t_name
        MERGE (a)-[:`${r}`]->(b)
        """, h=h, t=t, h_name=h.split(":")[1], t_name=t.split(":")[1])

# 4. 查询
with driver.session() as s:
    result = s.run("""
    MATCH (c:Component {name: 'compA'})-[:MAY_HAVE_FAULT]->(f:Fault)
    RETURN f.name AS fault
    """)
    for row in result:
        print("Fault:", row["fault"])

driver.close()

Logo

这里是“一人公司”的成长家园。我们提供从产品曝光、技术变现到法律财税的全栈内容,并连接云服务、办公空间等稀缺资源,助你专注创造,无忧运营。

更多推荐