边缘计算与工业物联网(Modbus/OPC UA)
边缘计算架构
边缘节点层级
工业边缘计算采用三层物理架构,将算力从云端下沉至生产现场。
+----------------------------------------------------------------------+
| 云端层(Cloud Layer) |
| 工业云平台 | 数据中台 | AI 训练 | 远程监控 | 应用管理 |
+----------------------------------------------------------------------+
| WAN / 4G/5G
+----------------------------------------------------------------------+
| 边缘层(Edge Layer) |
| 边缘节点:协议转换 | 数据清洗 | 实时决策 | 本地缓存 | 断点续传 |
+----------------------------------------------------------------------+
| Ethernet / RS-485 / CAN
+----------------------------------------------------------------------+
| 现场层(Field Layer) |
| PLC | 电表 | 传感器 | 变频器 | 工业机器人 | 视觉相机 |
+----------------------------------------------------------------------+现场层:包含各类工业设备与传感器,通过 Modbus RTU/ASCII、PROFINET、EtherNet/IP 等工业协议与上层通信。
边缘层:部署边缘计算节点,承担协议解析、数据预处理、实时控制、本地推理等任务,是边缘计算的核心。
云端层:提供海量数据存储、全局模型训练、远程运维、应用编排等能力。
边缘计算核心价值
| 维度 | 说明 |
|---|---|
| 低延迟 | 数据处理在本地完成,端到端延迟可控制在 10ms 以内,满足工业实时控制需求 |
| 带宽节省 | 仅上传聚合后的特征数据或告警信息,原始高频数据本地处理,可减少 90% 以上的上行带宽 |
| 本地自治 | 断网情况下边缘节点仍能独立运行,网络恢复后自动同步数据 |
| 数据安全 | 敏感工业数据不出厂区,仅上传脱敏或聚合结果,降低数据泄露风险 |
三类边缘节点
轻量网关:基于 ARM Cortex-A 系列处理器(如 Raspberry Pi、Rockchip RK3568),运行 Linux,适用于少量设备接入(< 50 点)、数据转发与协议转换场景。成本低,功耗 5-15W。
智能网关:基于 x86 或高性能 ARM 处理器(如 Intel Atom、NVIDIA Jetson),支持容器化部署,可运行 KubeEdge / EdgeX Foundry 等框架,适用于中等规模产线(50-2000 点),具备本地推理能力。
边缘服务器:基于 x86 服务器(如 Intel Xeon、AMD EPYC),部署于车间级机房,支持虚拟化与 GPU 加速,可管理数千个设备点,运行完整工业 PaaS 平台。
边缘计算框架
KubeEdge
KubeEdge 将 Kubernetes 从云端扩展到边缘,实现云边协同的容器化管理。
核心组件:
+-------------------+ +-------------------+
| Cloud Core | | Edge Core |
| (云端) | | (边缘节点) |
| | WebSocket | |
| - CloudHub |<--------->| - EdgeHub |
| - DeviceController| gRPC | - Edged (轻量kubelet)|
| - SyncController | | - DeviceTwin |
| - Router | | - EventBus |
| | | - MetaManager |
+-------------------+ +-------------------+- CloudHub:负责云端与边缘节点的 WebSocket 连接管理,支持多租户隔离。
- EdgeHub:边缘端通信模块,与 CloudHub 保持双向同步,支持离线缓存。
- Edged:精简版 kubelet,管理边缘节点上的 Pod 生命周期,资源占用低。
- DeviceTwin:设备数字孪生,维护设备状态的双向同步,离线时本地存储设备属性。
- EventBus:基于 MQTT 协议与设备通信,支持与 EdgeX Foundry 集成。
云边协同流程:
KubeEdge 设备管理流程:
1. 云端创建 Device CRD → CloudHub 同步至边缘
2. EdgeHub 接收 → DeviceTwin 更新本地状态
3. DeviceTwin 通过 EventBus 下发配置到物理设备
4. 设备上报数据 → EventBus → DeviceTwin → EdgeHub → CloudHub → 云端更新状态EdgeX Foundry
EdgeX Foundry 是由 Linux 基金会管理的边缘计算微服务平台,采用松耦合微服务架构。
微服务架构:
+------------------------------------------------------------------+
| EdgeX Foundry |
| |
| +----------------+ +----------------+ +------------------+ |
| | Device Services| | Core Services | | Supporting | |
| | | | | | Services | |
| | - Modbus | | - Core Data | | - Rules Engine | |
| | - OPC UA | | - Metadata | | - Scheduler | |
| | - BACnet | | - Command | | - Alerts | |
| | - MQTT | | | | | |
| +----------------+ +----------------+ +------------------+ |
| |
| +----------------+ +-------------------------------------+ |
| | Export Services| | System Management | |
| | - HTTP/MQTT | | - Config / Logging / Metrics | |
| | - App Service | | | |
| +----------------+ +-------------------------------------+ |
+------------------------------------------------------------------+设备服务(Device Services):直接与物理设备交互的协议适配层。EdgeX 提供 Modbus、OPC UA、MQTT、BACnet、SNMP 等多个预建设备服务,也支持自定义 Device Service。
核心服务(Core Services):
- Core Data:持久化设备采集数据,提供 REST API 查询,支持数据清理策略。
- Metadata:管理设备、设备配置文件(Device Profile)、设备服务等元数据。
- Command:提供向设备下发命令的统一接口,屏蔽底层协议差异。
规则引擎(Rules Engine):基于 eKuiper(LF Edge 项目)实现 SQL 驱动的流式规则处理,支持在边缘端进行实时过滤、聚合、告警。
设备接入示例(Modbus Device Service):
# device-modbus/profile/temperature-sensor.yaml
name: "TemperatureSensor"
manufacturer: "IndustrialCo"
model: "IT-2000"
labels:
- "temperature"
- "modbus"
deviceResources:
- name: "temperature"
description: "当前温度值"
attributes:
{ primaryTable: "HOLDING_REGISTERS", startingAddress: 0, count: 1 }
properties:
valueType: "Int16"
readWrite: "R"
units: "Celsius"
scale: 0.1
- name: "humidity"
description: "当前湿度值"
attributes:
{ primaryTable: "HOLDING_REGISTERS", startingAddress: 1, count: 1 }
properties:
valueType: "Int16"
readWrite: "R"
units: "Percent"
scale: 0.1
deviceCommands:
- name: "readAll"
readWrite: "R"
resourceOperations:
- { deviceResource: "temperature" }
- { deviceResource: "humidity" }OpenYurt
OpenYurt(原 CoreRAS)是阿里云开源的边缘容器平台,通过 YurtHub 组件实现边缘节点与云端控制面的无缝连接。
核心特性:
- YurtHub:边缘端流量代理,缓存云端 API 响应,支持边缘节点离线自治。
- YurtTunnel:反向通道,支持云端访问边缘端服务(如 kubectl exec/logs)。
- YurtController:管理节点池(NodePool),实现节点单元化管理。
- YurtAppManager:支持边缘单元化部署(YurtAppSet),将应用按节点池分发。
节点池管理:
apiVersion: apps.openyurt.io/v1beta1
kind: NodePool
metadata:
name: beijing-factory
spec:
type: Edge
labels:
region: beijing
factory: plant-a
---
apiVersion: apps.openyurt.io/v1beta1
kind: YurtAppSet
metadata:
name: edge-inference
spec:
selector:
matchLabels:
app: inference
workloadTemplates:
- template:
replicas: 1
template:
spec:
containers:
- name: inference
image: registry/edge-inference:1.0
nodepoolSelector:
matchLabels:
region: beijing框架对比表
| 维度 | KubeEdge | EdgeX Foundry | OpenYurt |
|---|---|---|---|
| 定位 | 云边协同 K8s 扩展 | 边缘 IoT 微服务平台 | 边缘容器平台 |
| 基础 | Kubernetes | 原生微服务 | Kubernetes |
| 设备管理 | Device CRD + DeviceTwin | 完整 Device Service 生态 | 需集成三方方案 |
| 协议适配 | 需集成 EdgeX 或 MQTT | 内置 Modbus/OPC UA/BACnet 等 | 无内置,需自行适配 |
| 离线自治 | 支持(EdgeHub 缓存) | 支持(本地存储) | 支持(YurtHub 缓存) |
| AI 推理 | 支持 GPU/NPU 调度 | 通过 App Service 集成 | 支持 GPU 调度 |
| 适用场景 | 云边协同容器化部署 | 设备接入与数据采集 | K8s 边缘化改造 |
| 社区活跃度 | CNCF 毕业项目 | LF Edge 项目 | CNCF 沙箱项目 |
边缘 AI
模型压缩与剪枝
边缘设备算力与内存有限,需要对深度学习模型进行压缩才能部署。
剪枝(Pruning):移除网络中不重要的权重或神经元连接。
结构化剪枝流程:
原始模型 → 训练 → 评估权重重要性 → 移除低权重通道/层 → 微调恢复精度
↓
缩减体积的模型量化(Quantization):将 FP32 权重转换为 INT8 或 FP16,减少模型体积与推理延迟。
| 量化方式 | 精度损失 | 推理加速 | 适用场景 |
|---|---|---|---|
| FP32 -> FP16 | < 0.1% | 1.5-2x | GPU/NPU 推理 |
| FP32 -> INT8 | 0.5-2% | 3-4x | CPU/GPU 推理 |
| FP32 -> INT4 | 3-5% | 5-6x | 专用 NPU 推理 |
知识蒸馏(Knowledge Distillation):用一个大型教师模型指导一个小型学生模型训练,使学生模型在保持较小体积的同时近似教师模型的精度。
TensorRT / ONNX Runtime 推理引擎
TensorRT:NVIDIA 推出的 GPU 推理优化引擎,支持层融合、精度校准、动态张量内存管理。
import tensorrt as trt
# 构建 TensorRT 引擎
TRT_LOGGER = trt.Logger(trt.Logger.WARNING)
builder = trt.Builder(TRT_LOGGER)
network = builder.create_network(1 << int(trt.NetworkDefinitionCreationFlag.EXPLICIT_BATCH))
parser = trt.OnnxParser(network, TRT_LOGGER)
# 加载 ONNX 模型
with open("model.onnx", "rb") as f:
parser.parse(f.read())
# 配置推理精度(FP16 / INT8)
config = builder.create_builder_config()
config.set_flag(trt.BuilderFlag.FP16)
# 构建引擎并序列化
serialized_engine = builder.build_serialized_network(network, config)
with open("model.trt", "wb") as f:
f.write(serialized_engine)ONNX Runtime:微软开源的跨平台推理引擎,支持 CPU/GPU/NPU 多种后端,提供统一的推理接口。
import onnxruntime as ort
import numpy as np
# 创建 ONNX Runtime 推理会话
sess = ort.InferenceSession(
"model.onnx",
providers=["TensorrtExecutionProvider", "CUDAExecutionProvider", "CPUExecutionProvider"]
)
# 准备输入数据
input_name = sess.get_inputs()[0].name
input_data = np.random.randn(1, 3, 224, 224).astype(np.float32)
# 执行推理
outputs = sess.run(None, {input_name: input_data})
print(f"推理结果: {outputs[0].shape}")边缘推理部署流程
模型部署到边缘节点流程图
云端训练 模型优化 边缘部署
+----------+ +-------------+ +---------------+
| 训练数据 | | FP32 模型 | | 边缘推理服务 |
| 深度学习 | ONNX | 量化剪枝 | ONNX | GPU/NPU 加速 |
| 训练框架 | -------> | TensorRT | ------> | 本地推理缓存 |
| PyTorch/ | 导出 | 优化引擎 | 部署 | 结果上传云端 |
| TF/Paddle | | 精度验证 | | 模型热更新 |
+----------+ +-------------+ +---------------+离线场景下本地推理
边缘节点在断网时需自行完成推理,不依赖云端。
class EdgeInferenceEngine:
"""边缘推理引擎,支持离线模式"""
def __init__(self, model_path: str, cache_size: int = 1000):
self.session = ort.InferenceSession(
model_path,
providers=["CPUExecutionProvider"]
)
self.input_name = self.session.get_inputs()[0].name
self.output_name = self.session.get_outputs()[0].name
self.result_cache = [] # 推理结果本地缓存
self.cache_size = cache_size
self.online = True # 网络状态标记
def predict(self, sensor_data: np.ndarray) -> dict:
"""执行本地推理"""
# 数据预处理
input_tensor = self._preprocess(sensor_data)
# 模型推理
output = self.session.run(
[self.output_name], {self.input_name: input_tensor}
)[0]
result = self._postprocess(output)
# 缓存推理结果
self.result_cache.append({
"timestamp": time.time(),
"input_shape": sensor_data.shape,
"output": result.tolist()
})
if len(self.result_cache) > self.cache_size:
self.result_cache.pop(0)
return {
"prediction": result,
"source": "edge_local",
"online": self.online
}
def sync_to_cloud(self):
"""网络恢复后同步缓存结果到云端"""
if self.online and self.result_cache:
# 批量上传缓存数据
batch_data = self.result_cache[:]
self._upload_batch(batch_data)
self.result_cache.clear()
def _preprocess(self, data: np.ndarray) -> np.ndarray:
# 标准化、归一化等预处理
mean = np.array([0.485, 0.456, 0.406])
std = np.array([0.229, 0.224, 0.225])
return (data / 255.0 - mean) / std
def _postprocess(self, output: np.ndarray) -> np.ndarray:
# 后处理,如 softmax 阈值过滤
return np.argmax(output, axis=1)
def _upload_batch(self, batch: list):
"""上传缓存数据到云端 API"""
import requests
try:
requests.post(
"https://cloud-api.factory/edge/sync",
json={"results": batch},
timeout=5
)
except requests.RequestException:
self.online = FalseModbus 协议
Modbus 是工业自动化领域最广泛使用的串行通信协议,由 Modicon 公司于 1979 年提出,现由 Modbus 组织维护。
Modbus RTU / ASCII / TCP 三种模式
| 模式 | 传输介质 | 帧格式 | 数据编码 | 典型速率 | 校验方式 |
|---|---|---|---|---|---|
| Modbus RTU | RS-232 / RS-485 | 二进制 | 8 位数据 | 9600-115200 bps | CRC-16 |
| Modbus ASCII | RS-232 / RS-485 | ASCII 字符 | 7 位数据 | 9600-19200 bps | LRC |
| Modbus TCP | 以太网 (TCP/IP) | 二进制 | 8 位数据 | 10/100 Mbps | 无(依赖 TCP) |
Modbus RTU 帧结构:
+---------+-----------+-------------+---------+
| 地址域 | 功能码 | 数据域 | CRC-16 |
| 1 字节 | 1 字节 | 0-252 字节 | 2 字节 |
+---------+-----------+-------------+---------+
示例(读取保持寄存器 01H ~ 02H,从站地址 0x01):
请求:01 03 00 00 00 02 C4 0B
响应:01 03 04 00 64 00 32 F8 17Modbus TCP 帧结构:
+-----------+---------+-----------+-------------+---------+
| 事务标识符 | 协议标识 | 长度字段 | 单元标识符 | PDU |
| 2 字节 | 2 字节 | 2 字节 | 1 字节 | 可变 |
+-----------+---------+-----------+-------------+---------+
MBAP Header (7 字节) + PDU(功能码 + 数据)Modbus 数据模型
Modbus 协议定义四种基本数据对象,每种对象可通过不同的功能码进行访问。
| 数据对象 | 数据大小 | 读写属性 | 地址范围 | 说明 |
|---|---|---|---|---|
| 线圈(Coils) | 1 bit | 读写 | 00001-09999 | 数字量输出,如继电器 |
| 离散输入(Discrete Inputs) | 1 bit | 只读 | 10001-19999 | 数字量输入,如限位开关 |
| 保持寄存器(Holding Registers) | 16 bit | 读写 | 40001-49999 | 模拟量输出/参数,如设定值 |
| 输入寄存器(Input Registers) | 16 bit | 只读 | 30001-39999 | 模拟量输入,如传感器读数 |
功能码
| 功能码 | 名称 | 说明 |
|---|---|---|
| 01 (0x01) | Read Coils | 读取线圈状态 |
| 02 (0x02) | Read Discrete Inputs | 读取离散输入状态 |
| 03 (0x03) | Read Holding Registers | 读取保持寄存器 |
| 04 (0x04) | Read Input Registers | 读取输入寄存器 |
| 05 (0x05) | Write Single Coil | 写单个线圈 |
| 06 (0x06) | Write Single Register | 写单个寄存器 |
| 15 (0x0F) | Write Multiple Coils | 写多个线圈 |
| 16 (0x10) | Write Multiple Registers | 写多个寄存器 |
| 23 (0x17) | Read/Write Multiple Registers | 同时读写多个寄存器 |
地址映射配置
工业边缘网关通常通过 YAML 或 JSON 配置设备地址映射关系:
# modbus-device-mapping.yaml
devices:
- name: "PLC_ProductionLine_1"
protocol: "modbus-tcp"
host: "192.168.1.100"
port: 502
slave_id: 1
poll_interval: 1000 # ms
points:
- name: "machine_status"
address: 0 # 数据类型: 线圈,地址 00001
type: "coil"
description: "设备运行状态 (0=停止, 1=运行)"
- name: "temperature"
address: 0 # 数据类型: 保持寄存器,地址 40001
type: "holding_register"
data_type: "int16"
scale: 0.1
unit: "Celsius"
- name: "pressure"
address: 1 # 数据类型: 保持寄存器,地址 40002
type: "holding_register"
data_type: "float32"
unit: "MPa"
- name: "PowerMeter_WorkshopA"
protocol: "modbus-rtu"
serial_port: "COM3"
baud_rate: 9600
data_bits: 8
stop_bits: 1
parity: "even"
slave_id: 10
poll_interval: 2000
points:
- name: "voltage"
address: 0 # 数据类型: 输入寄存器,地址 30001
type: "input_register"
data_type: "int16"
scale: 0.1
unit: "V"
- name: "current"
address: 1
type: "input_register"
data_type: "int16"
scale: 0.01
unit: "A"
- name: "power"
address: 2
type: "input_register"
data_type: "uint32"
unit: "kW"常见设备接入示例
PLC 数据采集(Modbus TCP):
import struct
from pymodbus.client import ModbusTcpClient
class PLCCollector:
"""西门子 S7-1200 PLC 数据采集"""
def __init__(self, host: str, port: int = 502, unit: int = 1):
self.client = ModbusTcpClient(host, port=port)
self.unit = unit
def collect_temperature(self) -> float:
"""读取温度值(保持寄存器 40001)"""
response = self.client.read_holding_registers(
address=0, count=1, slave=self.unit
)
if response.isError():
raise RuntimeError(f"Modbus 读取错误: {response}")
# 原始值为 int16,按比例换算
raw_value = response.registers[0]
return struct.unpack('>h', struct.pack('>H', raw_value))[0] * 0.1
def get_machine_status(self) -> bool:
"""读取设备运行状态(线圈 00001)"""
response = self.client.read_coils(address=0, count=1, slave=self.unit)
if response.isError():
raise RuntimeError(f"Modbus 读取错误: {response}")
return response.bits[0]
def set_speed(self, speed: int):
"""设定电机转速(保持寄存器 40010)"""
response = self.client.write_register(address=9, value=speed, slave=self.unit)
if response.isError():
raise RuntimeError(f"Modbus 写入错误: {response}")
def close(self):
self.client.close()电表数据采集(Modbus RTU):
from pymodbus.client import ModbusSerialClient
class PowerMeterCollector:
"""多功能电表数据采集(Modbus RTU)"""
def __init__(self, port: str = "COM3", baudrate: int = 9600, unit: int = 10):
self.client = ModbusSerialClient(
port=port,
baudrate=baudrate,
bytesize=8,
parity="E",
stopbits=1,
timeout=3
)
def read_voltage(self) -> float:
"""读取电压(输入寄存器 30001,两个寄存器组成 Float32)"""
response = self.client.read_input_registers(
address=0, count=2, slave=10
)
if response.isError():
raise RuntimeError(f"电表读取错误: {response}")
# 大端序 IEEE 754 Float32
packed = struct.pack('>HH', response.registers[0], response.registers[1])
return struct.unpack('>f', packed)[0]
def read_energy_total(self) -> float:
"""读取总有功电能(输入寄存器 30003-30004)"""
response = self.client.read_input_registers(
address=2, count=2, slave=10
)
if response.isError():
raise RuntimeError(f"电表读取错误: {response}")
packed = struct.pack('>HH', response.registers[0], response.registers[1])
return struct.unpack('>f', packed)[0]
def close(self):
self.client.close()传感器数据采集(温度 + 湿度):
from pymodbus.client import ModbusTcpClient
class SensorHub:
"""工业传感器采集站,4 路传感器通过 RS-485 接入"""
def __init__(self, host: str = "192.168.1.200"):
self.client = ModbusTcpClient(host, port=502)
def read_all_sensors(self) -> dict:
"""批量读取 4 路传感器数据"""
results = {}
for sensor_id in range(1, 5):
response = self.client.read_input_registers(
address=0, count=2, slave=sensor_id
)
if response.isError():
results[f"sensor_{sensor_id}"] = {"error": str(response)}
continue
packed = struct.pack('>HH', response.registers[0], response.registers[1])
temperature, humidity = struct.unpack('>ff', packed)
results[f"sensor_{sensor_id}"] = {
"temperature": round(temperature, 2),
"humidity": round(humidity, 2),
"unit": {"temperature": "Celsius", "humidity": "Percent"}
}
return results
def close(self):
self.client.close()OPC UA
OPC UA(Unified Architecture,统一架构)是 OPC 基金会推出的工业通信标准,相比 Classic OPC(基于 COM/DCOM),OPC UA 具有平台无关、安全可靠、面向服务等特点。
OPC UA 架构
OPC UA 支持两种通信模式:
Client/Server 模式:
+----------------+ +----------------+
| OPC UA Client | 请求/响应 | OPC UA Server |
| | ===============> | |
| - 订阅数据变化 | 会话 | - 地址空间管理 |
| - 调用方法 | <=============== | - 订阅管理 |
| - 读写属性 | 通知 | - 历史数据存储 |
+----------------+ +----------------+PubSub 模式(OPC UA Part 14):
+----------------+ +----------------+
| Publisher | MQTT / UDP | Subscriber |
| (OPC UA Server)| ==============> | (OPC UA Client)|
| | 数据报文 | |
+----------------+ +----------------+PubSub 模式适用于大规模部署和云边协同场景,发布者与订阅者解耦,通过消息中间件传递数据。支持 JSON 和 UADP(UA Datagram Packet)两种编码格式。
信息模型
OPC UA 信息模型使用节点(Node)和引用(Reference)构建地址空间(Address Space)。
节点类型:
| 节点类型 | 说明 | 示例 |
|---|---|---|
| Object(对象) | 表示物理或逻辑实体 | 设备、控制器、产线 |
| Variable(变量) | 表示值,可读写或订阅 | 温度值、转速、状态 |
| Method(方法) | 可调用的操作 | 启动/停止设备、复位 |
| ObjectType(对象类型) | 对象的类型定义 | 泵类型、电机类型 |
| VariableType(变量类型) | 变量的类型定义 | 模拟量类型、枚举类型 |
| DataType(数据类型) | 数据类型定义 | 结构体、枚举 |
节点结构:
NodeId: ns=2;s=TemperatureSensor1
- NodeClass: Variable
- BrowseName: TemperatureSensor1
- DisplayName: 温度传感器 1
- Description: 车间 A 线温度传感器
- Value: 25.3
- DataType: Double
- AccessLevel: CurrentRead
- UserAccessLevel: CurrentRead信息模型示例(使用 OPC UA 定义的设备信息模型):
<!-- OPC UA 节点定义 -->
<UANode i="1001" NodeClass="Variable">
<DisplayName>Pressure</DisplayName>
<Description>管道压力</Description>
<DataType i="11"> <!-- Double -->
</DataType>
<Value>
<uax:Double>0.85</uax:Double>
</Value>
<AccessLevel>3</AccessLevel> <!-- CurrentRead | CurrentWrite -->
</UANode>
<UANode i="1002" NodeClass="Method">
<DisplayName>Calibrate</DisplayName>
<Description>校准传感器</Description>
<MethodDeclarationId>
<Identifier>1002</Identifier>
</MethodDeclarationId>
</UANode>UA 二进制协议
OPC UA 二进制协议(UA Binary)是 OPC UA 的高效编码格式,使用 ASN.1 风格的序列化规则,在性能敏感场景下优先使用。
OPC UA 二进制消息帧结构:
+------------+------------+------------------+-----------+
| MessageType| ChunkType | MessageSize | Body |
| 3 字节 | 1 字节 | 4 字节 | 可变长度 |
+------------+------------+------------------+-----------+
MessageType: "MSG" / "OPN" / "CLO"
ChunkType: 'A' (Final), 'C' (Intermediate), 'E' (Error)
Body: 序列化的 UA 消息结构(使用 OPC UA Binary Encoding)对称加密消息头(SecureMessageHeader)包含安全协议版本、SecurityTokenID、SequenceNumber 和 RequestID,保证消息的完整性和顺序。
安全模型
OPC UA 安全模型基于三大安全策略:
| 安全策略 | 加密 | 签名 | 说明 |
|---|---|---|---|
| None | 无 | 无 | 仅适用于封闭网络,不推荐 |
| Basic128Rsa15 | AES-128 | RSA-SHA1 | 基本加密级别 |
| Basic256 | AES-256 | RSA-SHA1 | 强加密级别 |
| Basic256Sha256 | AES-256 | RSA-SHA256 | 最高安全级别 |
| Aes128Sha256RsaOaep | AES-128 | RSA-OAEP | 2020 新增,推荐 |
| Aes256Sha256RsaPss | AES-256 | RSA-PSS | 最高级,推荐 |
证书管理:
OPC UA 使用 X.509 证书进行身份验证:
证书管理流程:
服务器启动 → 加载证书 → 检查有效期 → 建立安全通道
↓
拒绝未授权客户端
客户端连接 → 发送证书 → 验证服务器证书 → 协商安全策略 → 建立会话应用证书验证:
from opcua import Client
from opcua.crypto.security_policies import SecurityPolicyBasic256Sha256
class SecureOPCUAClient:
"""使用签名的 OPC UA 客户端"""
def __init__(self, endpoint: str):
self.client = Client(endpoint)
# 设置客户端证书和私钥
self.client.set_security(
SecurityPolicyBasic256Sha256,
certificate="certs/client_cert.der",
private_key="certs/client_key.pem",
server_certificate="certs/server_cert.der"
)
def connect_and_read(self, node_id: str):
try:
self.client.connect()
node = self.client.get_node(node_id)
value = node.get_value()
return value
finally:
self.client.disconnect()OPC UA 与 Modbus 对比
| 维度 | OPC UA | Modbus |
|---|---|---|
| 通信模式 | Client/Server + PubSub 发布订阅 | Master/Slave 主从模式 |
| 网络层 | TCP/IP, UDP (PubSub) | RS-485/RS-232 串口, TCP/IP |
| 数据模型 | 丰富的信息模型(对象/变量/方法/事件) | 简单的寄存器/线圈模型 |
| 发现机制 | 支持 LDS-ME(本地发现服务器),mDNS | 无,需静态配置设备地址 |
| 安全性 | 内建:签名 + 加密 + 证书管理 | 无内建安全机制,需链路层或 VPN 保护 |
| 互操作性 | 强,OPC UA 规范严格,信息模型标准化 | 弱,设备地址映射依赖厂商自定义 |
| 实时性 | 一般,适用于监控与数据采集 | 好,适用于实时控制 |
| 复杂度 | 高,学习曲线陡峭 | 低,协议简单易用 |
| 适用场景 | 跨平台系统集成、复杂信息模型、安全敏感场景 | 简单设备接入、实时控制、资源受限设备 |
工业数据采集
采集任务配置
工业数据采集需定义采集点、采集周期、协议类型和解析脚本,通过统一的采集任务配置管理。
# data-collection-task.yaml
collection_task:
name: "workshop_a_collection"
version: "2.1.0"
enabled: true
schedule:
type: "periodic" # periodic | on_change | on_demand
interval_ms: 1000 # 1 秒采集周期
cron: "" # 也支持 cron 表达式
jitter_ms: 100 # 随机抖动,防止设备端同时并发
data_sources:
- name: "plc_line_1"
protocol: "modbus-tcp"
host: "192.168.1.100"
port: 502
slave_id: 1
retry:
max_retries: 3
retry_delay_ms: 500
timeout_ms: 3000
points:
- id: "temp_01"
name: "炉温"
address: { table: "HOLDING_REGISTERS", start: 0, count: 1 }
data_type: "int16"
scale: 0.1
unit: "Celsius"
quality_enabled: true
- name: "sensor_array_1"
protocol: "opcua"
endpoint: "opc.tcp://192.168.1.200:4840"
security_policy: "Basic256Sha256"
security_mode: "SignAndEncrypt"
certificate: "certs/opcua_client.der"
retry:
max_retries: 2
retry_delay_ms: 1000
points:
- id: "pressure_01"
node_id: "ns=2;s=PressureSensor1"
data_type: "float64"
unit: "MPa"
parsing_scripts:
- name: "custom_parser_1"
language: "lua"
script: |
-- 自定义采集值解析脚本
function parse(raw_value, point_config)
-- 根据设备手册自定义解析逻辑
if point_config.data_type == "int16" then
local value = raw_value * point_config.scale
if value > 1000 then
-- 异常值标记处理
return { value = 0, quality = "Bad", reason = "value_out_of_range" }
end
return { value = value, quality = "Good" }
end
return { value = raw_value, quality = "Good" }
end数据质量标记
工业数据采集必须标记每条数据的质量,以便上层应用判断数据是否可用。
| 质量标记 | 类型 | 说明 |
|---|---|---|
| Good | 正常 | 数据采集正常,值可信 |
| Good (LocalOverride) | 正常 | 值为本地手动设定 |
| Uncertain | 不确定 | 数据存在一定程度的不确定性 |
| Uncertain (SensorNotAccurate) | 不确定 | 传感器精度超限 |
| Uncertain (LastUsableValue) | 不确定 | 使用最后一次有效值 |
| Bad | 异常 | 数据不可用 |
| Bad (DeviceFailure) | 异常 | 设备故障 |
| Bad (SensorFailure) | 异常 | 传感器故障 |
| Bad (OutOfService) | 异常 | 设备离线/停止服务 |
| Bad (Timeout) | 异常 | 采集超时 |
| Bad (ConfigError) | 异常 | 采集配置错误 |
数据质量枚举实现:
public enum DataQuality
{
Good = 0x00,
Good_LocalOverride = 0x40,
Uncertain = 0x80,
Uncertain_SensorNotAccurate = 0x84,
Uncertain_LastUsableValue = 0x8A,
Bad = 0xC0,
Bad_DeviceFailure = 0xC4,
Bad_SensorFailure = 0xC8,
Bad_OutOfService = 0xCC,
Bad_Timeout = 0xD0,
Bad_ConfigError = 0xD4
}
public struct TagValue
{
public string TagId { get; set; }
public object Value { get; set; }
public DataQuality Quality { get; set; }
public DateTime Timestamp { get; set; }
public string Source { get; set; }
}断点续传与本地缓存
边缘节点在网络不稳定的场景下,需实现断点续传与本地缓存机制,确保数据不丢失。
本地缓存策略:
采集数据 → 内存队列(RingBuffer) → 持久化缓存(SQLite / 本地文件)
↓
网络恢复 → 批量上传云端
↓
上传成功 → 清理已同步数据
↓
上传失败 → 保留缓存,下次重试
缓存配置:
- 缓存上限:10 GB 或 7 天数据(以先到者为准)
- 缓存清理:FIFO + 优先级(重要数据优先保留)
- 压缩存储:数据按小时分片,Gzip 压缩断点续传实现:
import json
import sqlite3
import hashlib
from datetime import datetime, timedelta
from queue import Queue
from typing import Optional
class EdgeDataCache:
"""边缘数据缓存与断点续传"""
def __init__(self, db_path: str = "edge_cache.db", max_size_mb: int = 10240):
self.db_path = db_path
self.max_size_bytes = max_size_mb * 1024 * 1024
self._init_db()
def _init_db(self):
"""初始化 SQLite 缓存数据库"""
conn = sqlite3.connect(self.db_path)
conn.execute("""
CREATE TABLE IF NOT EXISTS data_cache (
id INTEGER PRIMARY KEY AUTOINCREMENT,
task_name TEXT NOT NULL,
point_id TEXT NOT NULL,
value TEXT NOT NULL,
quality TEXT NOT NULL,
timestamp TEXT NOT NULL,
checksum TEXT NOT NULL,
synced INTEGER DEFAULT 0,
created_at TEXT DEFAULT (datetime('now'))
)
""")
conn.execute("""
CREATE INDEX IF NOT EXISTS idx_synced
ON data_cache(synced, timestamp)
""")
conn.execute("""
CREATE INDEX IF NOT EXISTS idx_task
ON data_cache(task_name)
""")
conn.commit()
conn.close()
def write_data(self, task_name: str, point_id: str, value: float,
quality: str, timestamp: str):
"""写入一条采集数据到缓存"""
payload = json.dumps({"v": value, "q": quality, "t": timestamp})
checksum = hashlib.md5(payload.encode()).hexdigest()
conn = sqlite3.connect(self.db_path)
conn.execute(
"INSERT INTO data_cache (task_name, point_id, value, quality, timestamp, checksum) "
"VALUES (?, ?, ?, ?, ?, ?)",
(task_name, point_id, str(value), quality, timestamp, checksum)
)
conn.commit()
conn.close()
self._check_cache_size()
def get_unsynced_data(self, limit: int = 1000) -> list:
"""获取未同步的数据,用于断点续传"""
conn = sqlite3.connect(self.db_path)
cursor = conn.execute(
"SELECT id, task_name, point_id, value, quality, timestamp "
"FROM data_cache WHERE synced = 0 "
"ORDER BY id ASC LIMIT ?",
(limit,)
)
rows = cursor.fetchall()
conn.close()
return rows
def mark_synced(self, ids: list):
"""标记数据为已同步"""
if not ids:
return
conn = sqlite3.connect(self.db_path)
placeholders = ",".join("?" for _ in ids)
conn.execute(
f"UPDATE data_cache SET synced = 1 WHERE id IN ({placeholders})",
ids
)
conn.commit()
conn.close()
def _check_cache_size(self):
"""检查缓存总大小,超出上限时清理最旧数据"""
conn = sqlite3.connect(self.db_path)
cursor = conn.execute("SELECT SUM(LENGTH(value) + LENGTH(point_id)) FROM data_cache")
total = cursor.fetchone()[0] or 0
if total > self.max_size_bytes:
# 删除 20% 最旧的未同步数据(已同步数据优先清理)
conn.execute("""
DELETE FROM data_cache WHERE id IN (
SELECT id FROM data_cache
ORDER BY synced ASC, id ASC
LIMIT (SELECT COUNT(*) * 0.2 FROM data_cache)
)
""")
conn.commit()
conn.close()
def sync_to_cloud(self, cloud_api_url: str, batch_size: int = 500):
"""批量上传未同步数据到云端"""
unsynced = self.get_unsynced_data(limit=batch_size)
if not unsynced:
return True
import requests
batch = []
ids = []
for row in unsynced:
ids.append(row[0])
batch.append({
"task_name": row[1],
"point_id": row[2],
"value": float(row[3]),
"quality": row[4],
"timestamp": row[5]
})
try:
response = requests.post(
cloud_api_url,
json={"batch": batch},
headers={"Content-Type": "application/json"},
timeout=30
)
if response.status_code == 200:
self.mark_synced(ids)
return True
except requests.RequestException:
pass
return False时序数据库写入优化
工业数据采集产生大量时序数据,写入时序数据库时需采取批量优化策略。
写入优化策略:
- 批量写入:将多条数据打包一次写入,避免逐条插入的网络开销。
非优化方式:每条数据一次 HTTP 请求
时间 → |--INSERT--|--INSERT--|--INSERT--|--INSERT--|
网络往返 5ms × 4 = 20ms
批量优化方式:100 条数据一次请求
时间 → |------BATCH INSERT (100 records)------|
网络往返 5ms × 1 = 5ms + 序列化 1ms = 6ms
吞吐提升:约 17 倍数据压缩:对时间戳和标签列使用差值编码和字典编码,对数值列使用旋转门算法压缩。
分片写入:按时间分区(每天/每小时一个表或分区),避免单分区数据量过大。
异步写入:使用缓冲区队列,后台批量消费写入,不阻塞采集任务。
写入示例(InfluxDB):
from influxdb_client import InfluxDBClient, Point
from influxdb_client.client.write_api import SYNCHRONOUS, ASYNCHRONOUS
import datetime
import queue
import threading
class IndustrialTimeSeriesWriter:
"""工业时序数据写入器,支持批量写入和异步队列"""
def __init__(self, url: str, token: str, org: str, bucket: str,
batch_size: int = 500, flush_interval: int = 5):
self.client = InfluxDBClient(url=url, token=token, org=org)
self.bucket = bucket
self.batch_size = batch_size
self.flush_interval = flush_interval
self._buffer = []
self._lock = threading.Lock()
self._running = True
self._flush_thread = threading.Thread(target=self._periodic_flush, daemon=True)
self._flush_thread.start()
def write_point(self, measurement: str, tags: dict, fields: dict,
timestamp: Optional[datetime.datetime] = None):
"""写入单个采集点数据"""
point = Point(measurement)
for k, v in tags.items():
point.tag(k, v)
for k, v in fields.items():
point.field(k, v)
if timestamp:
point.time(timestamp)
else:
point.time(datetime.datetime.utcnow())
with self._lock:
self._buffer.append(point)
if len(self._buffer) >= self.batch_size:
self._flush()
def _flush(self):
"""将缓冲区数据批量写入 InfluxDB"""
if not self._buffer:
return
batch = self._buffer[:]
self._buffer.clear()
try:
write_api = self.client.write_api(write_options=ASYNCHRONOUS)
write_api.write(bucket=self.bucket, record=batch)
except Exception as e:
# 写入失败时重新入队,避免数据丢失
with self._lock:
self._buffer = batch + self._buffer
# 控制缓存上限,防止内存溢出
if len(self._buffer) > 100000:
self._buffer = self._buffer[-50000:]
print(f"时序数据库写入失败: {e}")
def _periodic_flush(self):
"""定时刷新缓冲区"""
import time
while self._running:
time.sleep(self.flush_interval)
with self._lock:
self._flush()
def close(self):
self._running = False
self._flush_thread.join(timeout=10)
with self._lock:
self._flush()
self.client.close()
# 使用示例
writer = IndustrialTimeSeriesWriter(
url="http://localhost:8086",
token="edge-token-123",
org="factory",
bucket="industrial_data"
)
# 批量写入温度传感器数据
for i in range(1000):
writer.write_point(
measurement="temperature",
tags={"sensor_id": "TS-001", "workshop": "A", "line": "L1"},
fields={"value": 25.3 + i * 0.01, "quality": 1}
)
writer.close()旋转门算法(用于数值压缩,减少写入量):
class SwingingDoorCompression:
"""旋转门压缩算法,用于时序数据趋势压缩"""
def __init__(self, compression_deviation: float):
self.deviation = compression_deviation # 压缩偏差阈值
self.last_stored_value: Optional[float] = None
self.last_stored_time: Optional[float] = None
self.upper_slope: float = float('-inf')
self.lower_slope: float = float('inf')
def compress(self, timestamp: float, value: float) -> Optional[tuple]:
"""
判断当前数据点是否需要存储。
返回 (timestamp, value) 如果需要存储,否则返回 None。
"""
if self.last_stored_value is None:
self.last_stored_value = value
self.last_stored_time = timestamp
return (timestamp, value)
if value == self.last_stored_value:
return None
dt = timestamp - self.last_stored_time
if dt <= 0:
return None
# 计算上下边界斜率
upper = (value + self.deviation - self.last_stored_value) / dt
lower = (value - self.deviation - self.last_stored_value) / dt
if upper < self.upper_slope:
self.upper_slope = upper
if lower > self.lower_slope:
self.lower_slope = lower
if self.upper_slope > self.lower_slope:
# 超出压缩门限,存储上一个点,重置
stored = (self.last_stored_time, self.last_stored_value)
self.last_stored_value = value
self.last_stored_time = timestamp
self.upper_slope = float('-inf')
self.lower_slope = float('inf')
return stored
return None
def flush(self) -> Optional[tuple]:
"""强制输出最后一个点"""
if self.last_stored_value is not None:
point = (self.last_stored_time, self.last_stored_value)
self.last_stored_value = None
self.last_stored_time = None
return point
return None