基于 OPC UA 的轻量化 PLC 实时监控平台 —— 工业现场数据“最后一公里”实践
“工业现场的痛苦,往往不在算法多高明,而在于数据‘拿不到、拿不全、拿不稳’。OPC UA 就是那个把 PLC 里的‘黑盒数据’变成‘透明资产’的标准接口。”
—— 哈尔滨工程大学《工业过程控制》课程核心思想延伸
一、实际应用场景描述
在现代智能制造、能源管控、楼宇自动化等场景中,PLC 作为现场控制核心,掌握着设备运行状态、工艺参数、报警信息等关键数据:
┌──────────────────────────────────────────────┐
│ 典型工业现场数据采集架构 │
│ │
│ [现场层] 设备与传感器 │
│ • 温度传感器、压力变送器、流量计 │
│ • 电机、阀门、变频器等执行机构 │
│ • 急停按钮、限位开关等安全器件 │
│ │ │
│ ▼ 硬接线 (IO/Link/Profinet) │
│ ┌────────────────────────────┐ │
│ │ PLC (S7-1200/1500, │ │
│ │ 欧姆龙NX, 三菱iQ-R) │ │
│ │ • 实时控制逻辑 │ │
│ │ • 数据寄存器(DB/D/M区) │ │
│ │ • 报警与诊断缓冲区 │ │
│ └────────────┬───────────────┘ │
│ │ 传统方式痛点 │
│ │ • 厂商私有协议 │
│ │ • DLL/ActiveX依赖 │
│ │ • 32位/64位兼容性问题 │
│ │ • 防火墙穿透困难 │
│ ▼ │
│ ┌────────────────────────────┐ │
│ │ OPC UA Server │ │
│ │ • 统一地址空间 │ │
│ │ • 内置安全机制(TLS/认证) │ │
│ │ • 跨平台、跨语言支持 │ │
│ │ • 历史数据+实时数据 │ │
│ └────────────┬───────────────┘ │
│ │ TCP/4840 (标准端口) │
│ ▼ │
│ ┌────────────────────────────┐ │
│ │ Python OPC UA Client │ ←── 本文重点 │
│ │ • 异步订阅(Subscription) │ │
│ │ • 断线重连机制 │ │
│ │ • 数据类型自动转换 │ │
│ │ • 轻量化监控平台后端 │ │
│ └────────────┬───────────────┘ │
│ │ REST API / WebSocket │
│ ▼ │
│ ┌────────────────────────────┐ │
│ │ 上位机监控平台 │ │
│ │ • Web 组态画面 │ │
│ │ • 实时趋势曲线 │ │
│ │ • 报警推送(钉钉/微信) │ │
│ │ • 移动端HMI │ │
│ └───────────────────────────┘ │
│ │
│ 核心价值: 一套代码,连接所有主流PLC │
│ 一次配置,打通OT与IT数据壁垒 │
└──────────────────────────────────────────────┘
典型应用场景
行业 应用场景 监控数据
离散制造 产线设备状态监控 伺服位置、IO状态、报警代码
流程工业 反应釜/锅炉监控 温度、压力、流量、液位
能源管理 光伏/风电/储能 电压、电流、功率、发电量
智慧楼宇 HVAC/照明/电梯 温湿度、CO₂、运行状态
水务环保 污水/净水处理 pH、溶解氧、浊度、加药量
物流仓储 AGV/堆垛机监控 位置、速度、电量、故障码
二、引入痛点
2.1 现场的真实困境
场景 现场发生了什么 根因
“驱动地狱” “西门子DLL只能在32位系统跑,服务器是64位” 厂商私有协议绑定
“协议迷宫” “S7、FINS、MC、Modbus…每种PLC一套代码” 缺乏统一标准
“数据孤岛” “PLC有数据,MES看不到,只能人工抄表” OT与IT系统割裂
“断线即崩” “网络抖动一下,采集程序就挂了” 缺乏健壮的重连机制
“安全隐患” “明文传输,工控网络暴露在公网” 传统OPC Classic安全性差
“性能瓶颈” “轮询1000个点,CPU占用90%+” 轮询机制效率低
2.2 核心矛盾
工业现场的核心矛盾是“控制系统的封闭性”与“数字化需求的开放性”之间的冲突。PLC 厂商倾向于通过私有协议锁定生态,而数字化转型需要开放、标准、安全的数据接口。OPC UA 正是为解决这一矛盾而生。
2.3 我们要解决什么
用一段精简的 Python 程序,构建一个 轻量级 OPC UA 监控客户端,实现:
1. 统一接口 —— 一套代码支持所有主流 PLC
2. 异步订阅 —— 基于事件驱动,而非低效轮询
3. 断线自愈 —— 自动重连,保证数据连续性
4. 类型安全 —— 自动处理 VARIANT 数据类型
5. 轻量部署 —— 无 DLL 依赖,跨平台运行
6. 快速上云 —— 为后续 MQTT/InfluxDB 预留接口
三、核心逻辑讲解
3.1 理论基础:OPC UA 通信模型
本工具基于哈工程《工业过程控制》第十三章“计算机过程控制系统”和 OPC UA 规范 Part 4(服务):
① OPC UA 核心概念
OPC UA 地址空间模型:
┌─────────────────────────────────────────────┐
│ RootFolder │
│ ├── Objects │
│ │ ├── Server │
│ │ │ ├── ServerStatus │
│ │ │ └── ServiceLevel │
│ │ └── DeviceSet │
│ │ ├── PLC_1 │
│ │ │ ├── DI (数据块) │
│ │ │ │ ├── Temperature (Double) │
│ │ │ │ ├── Pressure (Float) │
│ │ │ │ └── Motor_Status (Bool) │
│ │ │ └── AI (模拟量输入) │
│ │ └── PLC_2 │
│ ├── Types │
│ │ ├── ObjectTypes │
│ │ ├── VariableTypes │
│ │ └── DataTypes │
│ └── Views │
└─────────────────────────────────────────────┘
节点(Node)属性:
• NodeId: 唯一标识 (NamespaceIndex + Identifier)
• BrowseName: 浏览名 (不翻译)
• DisplayName: 显示名 (可本地化)
• Description: 描述信息
• Value: 当前值 (仅Variable节点)
• DataType: 数据类型 (Int32, Float, Double, Bool…)
• AccessLevel: 访问权限 (读/写)
• Historizing: 是否支持历史读取
② 订阅(Subscription)机制
OPC UA 发布-订阅模型:
[OPC UA Server] [OPC UA Client]
│ │
│ 1. 创建Session (CreateSession) │
│◄───────────────────────────────────────────│
│ 2. 激活Session (ActivateSession) │
│──────────────────────────────────────────►│
│ 3. 创建Subscription (CreateSubscription) │
│◄───────────────────────────────────────────│
│ 4. 创建MonitoredItem (CreateMonitoredItems)│
│──────────────────────────────────────────►│
│ │
│ 数据变化时… │
│ ┌──────────────────────────────────────┐│
│ │ PublishResponse (DataChange) ││
│ │ • ClientHandle ││
│ │ • Value (DataValue) ││
│ │ • StatusCode ││
│ │ • SourceTimestamp ││
│ │ • ServerTimestamp ││
│ └──────────────────────────────────────┘│
│──────────────────────────────────────────►│
│ │
│ 5. 定期保活 (PublishRequest) │
│◄───────────────────────────────────────────│
│ 6. 确认保活 (PublishResponse) │
│──────────────────────────────────────────►│
优势:
• 事件驱动:仅在数据变化时推送,节省带宽
• 批量传输:多个变化合并在一个消息中
• 断线续传:Server缓存未确认的Notification
• 服务质量:可配置发布间隔、队列大小、丢弃策略
③ 安全模型
OPC UA 安全层级:
应用层安全 (Application Authentication):
• 证书认证 (X.509)
• 用户名/密码
• 匿名访问 (不推荐)
传输层安全 (Transport Security):
• TLS 1.2/1.3 加密
• 端口: 4840 (opc.tcp)
消息层安全 (Message Security):
• Sign (签名): 防止篡改
• SignAndEncrypt (签名+加密): 最高安全
会话层安全 (Session Security):
• Session密钥轮换
• 超时自动注销
3.2 监控平台架构设计
┌──────────────────────────────────────────────────────┐
│ 轻量化监控平台架构 │
│ │
│ ┌────────────────────────────────────────────────┐ │
│ │ OPC UA Client Core │ │
│ │ ┌──────────────────────────────────────────┐ │ │
│ │ │ ConnectionManager (连接管理器) │ │ │
│ │ │ • 自动重连机制 │ │ │
│ │ │ • 心跳检测 │ │ │
│ │ │ • 多Server负载均衡 │ │ │
│ │ └──────────────┬───────────────────────────┘ │ │
│ │ │ │ │
│ │ ┌──────────────▼───────────────────────────┐ │ │
│ │ │ SubscriptionManager (订阅管理器) │ │ │
│ │ │ • 创建/删除Subscription │ │ │
│ │ │ • 批量添加MonitoredItem │ │ │
│ │ │ • 发布间隔动态调整 │ │ │
│ │ └──────────────┬───────────────────────────┘ │ │
│ │ │ │ │
│ │ ┌──────────────▼───────────────────────────┐ │ │
│ │ │ DataHandler (数据处理回调) │ │ │
│ │ │ • 数据类型转换 │ │ │
│ │ │ • 质量戳校验 │ │ │
│ │ │ • 死区过滤 │ │ │
│ │ │ • 时间戳对齐 │ │ │
│ │ └──────────────┬───────────────────────────┘ │ │
│ │ │ │ │
│ │ ┌──────────────▼───────────────────────────┐ │ │
│ │ │ StorageAdapter (存储适配器) │ │ │
│ │ │ • 内存缓存 (RingBuffer) │ │ │
│ │ │ • SQLite (本地持久化) │ │ │
│ │ │ • InfluxDB (时序数据库) │ │ │
│ │ │ • MQTT (实时推送) │ │ │
│ │ └──────────────┬───────────────────────────┘ │ │
│ │ │ │ │
│ │ ┌──────────────▼───────────────────────────┐ │ │
│ │ │ AlarmEngine (报警引擎) │ │ │
│ │ │ • 阈值判断 │ │ │
│ │ │ • 变化率报警 │ │ │
│ │ │ • 报警抑制/延时 │ │ │
│ │ │ • 通知分发(Webhook/邮件) │ │ │
│ │ └──────────────────────────────────────────┘ │ │
│ └────────────────────────────────────────────────┘ │
│ │
│ ┌────────────────────────────────────────────────┐ │
│ │ Web API Layer │ │
│ │ • FastAPI / Flask │ │
│ │ • RESTful 接口 │ │
│ │ • WebSocket 实时推送 │ │
│ │ • JWT 身份认证 │ │
│ └────────────────────────────────────────────────┘ │
│ │
│ ┌────────────────────────────────────────────────┐ │
│ │ Frontend Layer │ │
│ │ • Vue/React 组态画面 │ │
│ │ • ECharts 趋势曲线 │ │
│ │ • Ant Design 控制面板 │ │
│ │ • 移动端自适应 │ │
│ └────────────────────────────────────────────────┘ │
│ │
│ 核心设计原则: │
│ • 高内聚低耦合: 每层只负责单一职责 │
│ • 异步非阻塞: 使用asyncio提高并发性能 │
│ • 优雅降级: 某模块故障不影响整体运行 │
│ • 可观测性: 内置日志、指标、追踪 │
└──────────────────────────────────────────────────────┘
四、代码讲解(面向对象设计)
4.1 类结构总览
类名 职责 设计模式
"OpcUaConfig" OPC UA 连接配置(dataclass) 值对象
"TagConfig" 标签配置(dataclass) 值对象
"DataPoint" 数据点(dataclass) 值对象
"ConnectionManager" 连接与重连管理 单例模式
"SubscriptionManager" 订阅生命周期管理 观察者模式
"DataCallbackHandler" 数据变更回调处理 策略模式
"StorageAdapter" 数据存储抽象接口 适配器模式
"MemoryStorage" 内存存储实现 具体实现
"AlarmEngine" 报警判断引擎 状态模式
"OpcUaMonitorClient" OPC UA 监控客户端(聚合根) 聚合根
"WebApiServer" Web API 服务 外观模式
4.2 核心代码(精简版,CSDN友好)
完整源码约 300 行,包含 8 个类、异步订阅、断线重连、内存存储、REST API。
以下为可直接运行的精简核心版。
"""
基于 OPC UA 的轻量化 PLC 实时监控平台
参考哈尔滨工程大学《工业过程控制》第十三章"计算机过程控制系统"
"""
from dataclasses import dataclass, field
from typing import List, Dict, Optional, Callable, Any, Set
from enum import Enum, auto
import asyncio
import logging
from datetime import datetime, timezone
import json
from collections import deque
import uuid
# OPC UA 依赖
try:
from asyncua import Client, ua
from asyncua.common.subscription import Subscription
from asyncua.common.node import Node
except ImportError:
print("请安装依赖: pip install asyncua")
exit(1)
# ============================================================
# 1. 基础数据结构(值对象)
# ============================================================
class DataQuality(Enum):
"""数据质量枚举"""
GOOD = auto()
UNCERTAIN = auto()
BAD = auto()
COMM_FAILURE = auto()
@dataclass
class OpcUaConfig:
"""OPC UA 连接配置 —— 值对象"""
server_url: str = "opc.tcp://localhost:4840"
security_mode: str = "None" # None, Sign, SignAndEncrypt
certificate_path: Optional[str] = None
private_key_path: Optional[str] = None
username: Optional[str] = None
password: Optional[str] = None
timeout: int = 5000 # ms
reconnect_interval: int = 5 # s
max_reconnect_attempts: int = 10
@dataclass
class TagConfig:
"""标签配置 —— 值对象"""
node_id: str # OPC UA NodeId, 如 "ns=3;s=DB1.Temperature"
display_name: str # 显示名称
data_type: type = float # Python数据类型
unit: str = "" # 单位
deadband: float = 0.0 # 死区(绝对值)
sampling_interval: int = 1000 # ms
queue_size: int = 10 # 队列大小
discard_oldest: bool = True # 丢弃旧数据策略
alarm_low: Optional[float] = None
alarm_high: Optional[float] = None
alarm_enabled: bool = False
@dataclass
class DataPoint:
"""数据点 —— 值对象"""
tag_name: str
value: Any
timestamp: datetime
quality: DataQuality
source_timestamp: Optional[datetime] = None
server_timestamp: Optional[datetime] = None
def to_dict(self) -> dict:
"""转换为字典,便于序列化"""
return {
"tag": self.tag_name,
"value": self.value,
"timestamp": self.timestamp.isoformat(),
"quality": self.quality.name,
"unit": ""
}
# ============================================================
# 2. 连接管理器(单例模式)
# ============================================================
class ConnectionManager:
"""OPC UA 连接管理器 —— 单例模式"""
_instance = None
_lock = asyncio.Lock()
def __new__(cls):
if cls._instance is None:
cls._instance = super().__new__(cls)
cls._instance._initialized = False
return cls._instance
def __init__(self):
if self._initialized:
return
self.config: Optional[OpcUaConfig] = None
self.client: Optional[Client] = None
self.connected = False
self.reconnect_task: Optional[asyncio.Task] = None
self.reconnect_count = 0
self.connection_lock = asyncio.Lock()
self.logger = self._setup_logger()
self._initialized = True
def _setup_logger(self) -> logging.Logger:
"""配置日志"""
logger = logging.getLogger(__name__)
if not logger.handlers:
handler = logging.StreamHandler()
formatter = logging.Formatter(
'%(asctime)s – %(name)s – %(levelname)s – %(message)s'
)
handler.setFormatter(formatter)
logger.addHandler(handler)
logger.setLevel(logging.INFO)
return logger
async def connect(self, config: OpcUaConfig) -> bool:
"""建立连接"""
async with self.connection_lock:
self.config = config
try:
# 创建客户端
self.client = Client(url=config.server_url, timeout=config.timeout)
# 配置安全策略
if config.security_mode == "Sign":
self.client.set_security_string("Basic256Sha256,Sign")
elif config.security_mode == "SignAndEncrypt":
self.client.set_security_string("Basic256Sha256,SignAndEncrypt")
# 配置认证
if config.username and config.password:
self.client.set_user(config.username)
self.client.set_password(config.password)
# 连接
await self.client.connect()
self.connected = True
self.reconnect_count = 0
self.logger.info(f"✅ 成功连接到 OPC UA Server: {config.server_url}")
# 获取Server状态
server_status = await self.client.get_node("i=2259").read_value()
self.logger.info(f" Server状态: {server_status}")
return True
except Exception as e:
self.logger.error(f"❌ 连接失败: {e}")
self.connected = False
return False
async def disconnect(self):
"""断开连接"""
async with self.connection_lock:
if self.client and self.connected:
try:
await self.client.disconnect()
self.logger.info("✅ 已断开 OPC UA 连接")
except Exception as e:
self.logger.error(f"❌ 断开连接时出错: {e}")
finally:
self.connected = False
self.client = None
async def ensure_connected(self) -> bool:
"""确保连接可用,必要时重连"""
if self.connected and self.client:
try:
# 心跳检测
await self.client.get_node("i=2259").read_value()
return True
except Exception:
self.logger.warning("⚠️ 连接已断开,准备重连…")
self.connected = False
if self.reconnect_count >= self.config.max_reconnect_attempts:
self.logger.error("❌ 达到最大重连次数,放弃重连")
return False
self.reconnect_count += 1
self.logger.info(f"🔄 尝试第 {self.reconnect_count} 次重连…")
await asyncio.sleep(self.config.reconnect_interval)
return await self.connect(self.config)
def start_reconnect_monitor(self):
"""启动重连监控(后台任务)"""
if self.reconnect_task is None or self.reconnect_task.done():
self.reconnect_task = asyncio.create_task(self._reconnect_monitor())
async def _reconnect_monitor(self):
"""重连监控循环"""
while True:
if not self.connected:
await self.ensure_connected()
await asyncio.sleep(self.config.reconnect_interval)
# ============================================================
# 3. 数据存储适配器(适配器模式)
# ============================================================
class StorageAdapter:
"""存储适配器抽象基类"""
async def store(self, data_point: DataPoint) -> bool:
raise NotImplementedError
async def query(self, tag_name: str, start_time: datetime,
end_time: datetime) -> List[DataPoint]:
raise NotImplementedError
async def close(self):
pass
class MemoryStorage(StorageAdapter):
"""内存存储实现 —— 轻量级方案"""
def __init__(self, max_points_per_tag: int = 10000):
self.data: Dict[str, deque] = {}
self.max_points_per_tag = max_points_per_tag
self.lock = asyncio.Lock()
async def store(self, data_point: DataPoint) -> bool:
async with self.lock:
if data_point.tag_name not in self.data:
self.data[data_point.tag_name] = deque(maxlen=self.max_points_per_tag)
self.data[data_point.tag_name].append(data_point)
return True
async def query(self, tag_name: str, start_time: datetime,
end_time: datetime) -> List[DataPoint]:
async with self.lock:
if tag_name not in self.data:
return []
result = []
for point in self.data[tag_name]:
if start_time <= point.timestamp <= end_time:
result.append(point)
return result
async def get_latest(self, tag_name: str) -> Optional[DataPoint]:
"""获取最新数据点"""
async with self.lock:
if tag_name in self.data and self.data[tag_name]:
return self.data[tag_name][-1]
return None
async def get_all_latest(self) -> Dict[str, DataPoint]:
"""获取所有标签的最新值"""
async with self.lock:
result = {}
for tag_name, points in self.data.items():
if points:
result[tag_name] = points[-1]
return result
# ============================================================
# 4. 报警引擎(状态模式)
# ============================================================
class AlarmState(Enum):
"""报警状态"""
NORMAL = auto()
LOW_ALARM = auto()
HIGH_ALARM = auto()
COMM_FAILURE = auto()
@dataclass
class AlarmEvent:
"""报警事件"""
tag_name: str
alarm_type: AlarmState
value: Any
threshold: Optional[float]
timestamp: datetime
message: str
class AlarmEngine:
"""报警引擎 —— 状态模式"""
def __init__(self):
self.alarm_sta
利用AI解决实际问题,如果你觉得这个工具好用,欢迎关注长安牧笛!




