知识图谱完整指南:概念、构建流程与 Python 实现
知识图谱完整指南:概念、构建流程与 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 管道监听数据库变更,实时更新图谱
- 多模态图谱: 将图片、运维记录也作为节点嵌入
八、总结
构建一个工业知识图谱的完整流程:
- 明确范围 → 选定领域、实体类型
- 设计本体 → 定义类和关系
- 数据接入 → 结构化直接转三元组,非结构化用 NLP 抽取
- 清洗融合 → 对齐、去重、消歧
- 存储查询 → Neo4j(大图)或 NetworkX(原型)
- 质量评估 → 专家验证、准确率指标
- 应用开发 → 问答、搜索、可视化、告警
附录:快速运行示例
# 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()
更多推荐



所有评论(0)