29.2 工业物联网架构与实现

工业物联网概述

工业物联网(Industrial Internet of Things, IIoT)是物联网技术在工业领域的应用,通过将传感器、执行器、控制器等设备连接到网络,实现设备间的数据交换和协同工作,为智能制造提供基础支撑。

工业物联网的核心特征:

  1. 设备互联

    • 传感器网络:温度、压力、振动、流量等传感器
    • 执行器控制:电机、阀门、机械臂等执行设备
    • 智能网关:协议转换、数据预处理、边缘计算
    • 通信协议:Modbus、OPC-UA、MQTT、CoAP等
  2. 数据采集

    • 实时数据:设备运行状态、生产参数
    • 历史数据:设备维护记录、生产历史
    • 事件数据:报警信息、异常事件
    • 视频数据:监控画面、质量检测图像
  3. 边缘处理

    • 本地计算:数据预处理、实时分析
    • 智能决策:基于规则的自动控制
    • 缓存机制:网络中断时的数据缓存
    • 安全防护:边缘安全、数据加密
  4. 云端集成

    • 数据上云:大数据存储和分析
    • 远程监控:设备状态远程查看
    • 预测维护:基于AI的故障预测
    • 系统集成:与ERP、MES等系统集成

工业物联网架构

工业物联网架构
感知层
网络层
边缘层
平台层
应用层
传感器
执行器
RFID/条码
工业相机
PLC/DCS
有线网络
无线网络
5G通信
LoRa/NB-IoT
边缘网关
边缘计算
本地存储
协议转换
IoT平台
数据管理
设备管理
安全管理
监控应用
分析应用
控制应用
维护应用
import numpy as np
import pandas as pd
import matplotlib.pyplot as plt
import seaborn as sns
from datetime import datetime, timedelta
from typing import Dict, List, Any, Optional, Tuple, Union
from dataclasses import dataclass, asdict, field
from enum import Enum
import json
import sqlite3
from pathlib import Path
import logging
from abc import ABC, abstractmethod
import random
import time
from concurrent.futures import ThreadPoolExecutor
import threading
from queue import Queue, Empty
import socket
import struct
import hashlib
import warnings
warnings.filterwarnings('ignore')

# 设置中文字体
plt.rcParams['font.sans-serif'] = ['SimHei', 'Arial Unicode MS']
plt.rcParams['axes.unicode_minus'] = False

class DeviceType(Enum):
    """设备类型枚举"""
    SENSOR = "sensor"
    ACTUATOR = "actuator"
    CONTROLLER = "controller"
    GATEWAY = "gateway"
    CAMERA = "camera"

class DataType(Enum):
    """数据类型枚举"""
    TEMPERATURE = "temperature"
    PRESSURE = "pressure"
    VIBRATION = "vibration"
    FLOW = "flow"
    VOLTAGE = "voltage"
    CURRENT = "current"
    SPEED = "speed"
    POSITION = "position"

class ConnectionStatus(Enum):
    """连接状态枚举"""
    ONLINE = "online"
    OFFLINE = "offline"
    ERROR = "error"
    MAINTENANCE = "maintenance"

@dataclass
class IoTDevice:
    """
    IoT设备信息
    """
    device_id: str
    device_name: str
    device_type: DeviceType
    location: str
    ip_address: str
    port: int
    protocol: str  # MQTT, Modbus, OPC-UA等
    data_types: List[DataType]
    sampling_rate: float  # 采样频率(Hz)
    status: ConnectionStatus = ConnectionStatus.OFFLINE
    last_seen: Optional[datetime] = None
    firmware_version: str = "1.0.0"
    battery_level: Optional[float] = None  # 电池电量(0-100)
    
    def __post_init__(self):
        if self.data_types is None:
            self.data_types = []

@dataclass
class SensorData:
    """
    传感器数据
    """
    device_id: str
    timestamp: datetime
    data_type: DataType
    value: float
    unit: str
    quality: float = 1.0  # 数据质量(0-1)
    location: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

@dataclass
class EdgeNode:
    """
    边缘节点信息
    """
    node_id: str
    node_name: str
    location: str
    ip_address: str
    cpu_cores: int
    memory_gb: float
    storage_gb: float
    connected_devices: List[str] = field(default_factory=list)
    processing_rules: List[Dict] = field(default_factory=list)
    status: ConnectionStatus = ConnectionStatus.OFFLINE
    cpu_usage: float = 0.0
    memory_usage: float = 0.0
    storage_usage: float = 0.0

class IIoTDataProcessor:
    """
    工业物联网数据处理器
    
    功能:
    - 数据采集和预处理
    - 实时数据分析
    - 异常检测
    - 数据质量评估
    """
    
    def __init__(self):
        """
        初始化数据处理器
        """
        self.processing_rules = []
        self.anomaly_thresholds = {}
        self.data_buffer = Queue(maxsize=10000)
        self.processed_data = []
        self.alerts = []
        
    def add_processing_rule(self, rule_name: str, condition: str, action: str, 
                          parameters: Dict[str, Any] = None):
        """
        添加数据处理规则
        
        Args:
            rule_name: 规则名称
            condition: 触发条件
            action: 执行动作
            parameters: 规则参数
        """
        rule = {
            'name': rule_name,
            'condition': condition,
            'action': action,
            'parameters': parameters or {},
            'created_at': datetime.now(),
            'enabled': True
        }
        self.processing_rules.append(rule)
    
    def set_anomaly_threshold(self, data_type: DataType, min_value: float, 
                            max_value: float, deviation_factor: float = 3.0):
        """
        设置异常检测阈值
        
        Args:
            data_type: 数据类型
            min_value: 最小正常值
            max_value: 最大正常值
            deviation_factor: 标准差倍数
        """
        self.anomaly_thresholds[data_type] = {
            'min_value': min_value,
            'max_value': max_value,
            'deviation_factor': deviation_factor
        }
    
    def process_sensor_data(self, sensor_data: SensorData) -> Dict[str, Any]:
        """
        处理传感器数据
        
        Args:
            sensor_data: 传感器数据
        
        Returns:
            处理结果
        """
        result = {
            'original_data': sensor_data,
            'processed_value': sensor_data.value,
            'quality_score': sensor_data.quality,
            'anomaly_detected': False,
            'alerts': [],
            'processed_at': datetime.now()
        }
        
        # 数据质量检查
        quality_issues = self._check_data_quality(sensor_data)
        if quality_issues:
            result['quality_score'] *= 0.8
            result['alerts'].extend(quality_issues)
        
        # 异常检测
        anomaly_result = self._detect_anomaly(sensor_data)
        if anomaly_result['is_anomaly']:
            result['anomaly_detected'] = True
            result['alerts'].append({
                'type': 'ANOMALY',
                'message': anomaly_result['message'],
                'severity': anomaly_result['severity']
            })
        
        # 应用处理规则
        rule_results = self._apply_processing_rules(sensor_data)
        result['rule_results'] = rule_results
        
        # 数据平滑和滤波
        result['processed_value'] = self._apply_smoothing(sensor_data)
        
        self.processed_data.append(result)
        return result
    
    def _check_data_quality(self, sensor_data: SensorData) -> List[Dict[str, str]]:
        """
        检查数据质量
        
        Args:
            sensor_data: 传感器数据
        
        Returns:
            质量问题列表
        """
        issues = []
        
        # 检查数据时效性
        time_diff = (datetime.now() - sensor_data.timestamp).total_seconds()
        if time_diff > 300:  # 5分钟
            issues.append({
                'type': 'TIMELINESS',
                'message': f'数据延迟{time_diff:.1f}秒',
                'severity': 'MEDIUM'
            })
        
        # 检查数值范围
        if sensor_data.data_type == DataType.TEMPERATURE:
            if sensor_data.value < -50 or sensor_data.value > 200:
                issues.append({
                    'type': 'RANGE',
                    'message': f'温度值{sensor_data.value}超出合理范围',
                    'severity': 'HIGH'
                })
        elif sensor_data.data_type == DataType.PRESSURE:
            if sensor_data.value < 0 or sensor_data.value > 1000:
                issues.append({
                    'type': 'RANGE',
                    'message': f'压力值{sensor_data.value}超出合理范围',
                    'severity': 'HIGH'
                })
        
        # 检查数据质量标识
        if sensor_data.quality < 0.8:
            issues.append({
                'type': 'QUALITY',
                'message': f'数据质量较低: {sensor_data.quality:.2f}',
                'severity': 'MEDIUM'
            })
        
        return issues
    
    def _detect_anomaly(self, sensor_data: SensorData) -> Dict[str, Any]:
        """
        检测数据异常
        
        Args:
            sensor_data: 传感器数据
        
        Returns:
            异常检测结果
        """
        result = {
            'is_anomaly': False,
            'message': '',
            'severity': 'LOW',
            'confidence': 0.0
        }
        
        if sensor_data.data_type not in self.anomaly_thresholds:
            return result
        
        threshold = self.anomaly_thresholds[sensor_data.data_type]
        value = sensor_data.value
        
        # 范围检查
        if value < threshold['min_value'] or value > threshold['max_value']:
            result['is_anomaly'] = True
            result['message'] = f'{sensor_data.data_type.value}值{value}超出正常范围[{threshold["min_value"]}, {threshold["max_value"]}]'
            result['severity'] = 'HIGH'
            result['confidence'] = 0.9
            return result
        
        # 获取历史数据进行统计分析
        recent_data = [d['original_data'].value for d in self.processed_data[-50:] 
                      if d['original_data'].data_type == sensor_data.data_type and 
                         d['original_data'].device_id == sensor_data.device_id]
        
        if len(recent_data) >= 10:
            mean_value = np.mean(recent_data)
            std_value = np.std(recent_data)
            
            # 3-sigma规则检测异常
            if abs(value - mean_value) > threshold['deviation_factor'] * std_value:
                result['is_anomaly'] = True
                result['message'] = f'{sensor_data.data_type.value}值{value}偏离均值{mean_value:.2f}超过{threshold["deviation_factor"]}倍标准差'
                result['severity'] = 'MEDIUM'
                result['confidence'] = min(0.8, abs(value - mean_value) / (threshold['deviation_factor'] * std_value))
        
        return result
    
    def _apply_processing_rules(self, sensor_data: SensorData) -> List[Dict[str, Any]]:
        """
        应用处理规则
        
        Args:
            sensor_data: 传感器数据
        
        Returns:
            规则执行结果
        """
        results = []
        
        for rule in self.processing_rules:
            if not rule['enabled']:
                continue
            
            try:
                # 简化的条件评估(实际应用中需要更复杂的规则引擎)
                condition_met = self._evaluate_condition(rule['condition'], sensor_data)
                
                if condition_met:
                    action_result = self._execute_action(rule['action'], sensor_data, rule['parameters'])
                    results.append({
                        'rule_name': rule['name'],
                        'condition': rule['condition'],
                        'action': rule['action'],
                        'result': action_result,
                        'executed_at': datetime.now()
                    })
            except Exception as e:
                results.append({
                    'rule_name': rule['name'],
                    'error': str(e),
                    'executed_at': datetime.now()
                })
        
        return results
    
    def _evaluate_condition(self, condition: str, sensor_data: SensorData) -> bool:
        """
        评估规则条件
        
        Args:
            condition: 条件表达式
            sensor_data: 传感器数据
        
        Returns:
            条件是否满足
        """
        # 简化的条件评估
        if 'temperature > 80' in condition and sensor_data.data_type == DataType.TEMPERATURE:
            return sensor_data.value > 80
        elif 'pressure < 1' in condition and sensor_data.data_type == DataType.PRESSURE:
            return sensor_data.value < 1
        elif 'vibration > 5' in condition and sensor_data.data_type == DataType.VIBRATION:
            return sensor_data.value > 5
        
        return False
    
    def _execute_action(self, action: str, sensor_data: SensorData, parameters: Dict[str, Any]) -> str:
        """
        执行规则动作
        
        Args:
            action: 动作类型
            sensor_data: 传感器数据
            parameters: 动作参数
        
        Returns:
            执行结果
        """
        if action == 'send_alert':
            alert = {
                'device_id': sensor_data.device_id,
                'data_type': sensor_data.data_type.value,
                'value': sensor_data.value,
                'message': parameters.get('message', '设备异常'),
                'severity': parameters.get('severity', 'MEDIUM'),
                'timestamp': datetime.now()
            }
            self.alerts.append(alert)
            return f"发送告警: {alert['message']}"
        
        elif action == 'log_event':
            return f"记录事件: {sensor_data.device_id} {sensor_data.data_type.value} = {sensor_data.value}"
        
        elif action == 'trigger_maintenance':
            return f"触发维护请求: {sensor_data.device_id}"
        
        return f"执行动作: {action}"
    
    def _apply_smoothing(self, sensor_data: SensorData) -> float:
        """
        应用数据平滑
        
        Args:
            sensor_data: 传感器数据
        
        Returns:
            平滑后的数值
        """
        # 获取同类型的最近数据
        recent_values = [d['original_data'].value for d in self.processed_data[-10:] 
                        if d['original_data'].data_type == sensor_data.data_type and 
                           d['original_data'].device_id == sensor_data.device_id]
        
        if len(recent_values) >= 3:
            # 使用移动平均进行平滑
            recent_values.append(sensor_data.value)
            return np.mean(recent_values[-5:])  # 5点移动平均
        
        return sensor_data.value
    
    def get_processing_statistics(self) -> Dict[str, Any]:
        """
        获取处理统计信息
        
        Returns:
            统计信息
        """
        if not self.processed_data:
            return {'message': '暂无处理数据'}
        
        total_processed = len(self.processed_data)
        anomaly_count = sum(1 for d in self.processed_data if d['anomaly_detected'])
        avg_quality = np.mean([d['quality_score'] for d in self.processed_data])
        
        # 按数据类型统计
        type_stats = {}
        for data in self.processed_data:
            data_type = data['original_data'].data_type.value
            if data_type not in type_stats:
                type_stats[data_type] = {'count': 0, 'anomalies': 0, 'avg_value': 0}
            
            type_stats[data_type]['count'] += 1
            if data['anomaly_detected']:
                type_stats[data_type]['anomalies'] += 1
        
        # 计算平均值
        for data_type in type_stats:
            values = [d['processed_value'] for d in self.processed_data 
                     if d['original_data'].data_type.value == data_type]
            type_stats[data_type]['avg_value'] = np.mean(values) if values else 0
        
        return {
            'total_processed': total_processed,
            'anomaly_count': anomaly_count,
            'anomaly_rate': anomaly_count / total_processed,
            'average_quality': avg_quality,
            'type_statistics': type_stats,
            'active_rules': len([r for r in self.processing_rules if r['enabled']]),
            'total_alerts': len(self.alerts)
        }

class EdgeComputingNode:
    """
    边缘计算节点
    
    功能:
    - 设备连接管理
    - 本地数据处理
    - 边缘AI推理
    - 云端数据同步
    """
    
    def __init__(self, node_info: EdgeNode):
        """
        初始化边缘计算节点
        
        Args:
            node_info: 节点信息
        """
        self.node_info = node_info
        self.connected_devices = {}
        self.data_processor = IIoTDataProcessor()
        self.local_storage = []
        self.sync_queue = Queue()
        self.is_running = False
        self.processing_thread = None
        
        # 初始化处理规则
        self._setup_default_rules()
        
        # 初始化异常检测阈值
        self._setup_anomaly_thresholds()
    
    def _setup_default_rules(self):
        """
        设置默认处理规则
        """
        # 温度过高告警
        self.data_processor.add_processing_rule(
            rule_name="高温告警",
            condition="temperature > 80",
            action="send_alert",
            parameters={'message': '设备温度过高', 'severity': 'HIGH'}
        )
        
        # 压力异常告警
        self.data_processor.add_processing_rule(
            rule_name="低压告警",
            condition="pressure < 1",
            action="send_alert",
            parameters={'message': '系统压力过低', 'severity': 'MEDIUM'}
        )
        
        # 振动异常维护
        self.data_processor.add_processing_rule(
            rule_name="振动异常",
            condition="vibration > 5",
            action="trigger_maintenance",
            parameters={'message': '设备振动异常,需要检查'}
        )
    
    def _setup_anomaly_thresholds(self):
        """
        设置异常检测阈值
        """
        self.data_processor.set_anomaly_threshold(DataType.TEMPERATURE, -10, 100, 3.0)
        self.data_processor.set_anomaly_threshold(DataType.PRESSURE, 0, 50, 2.5)
        self.data_processor.set_anomaly_threshold(DataType.VIBRATION, 0, 10, 2.0)
        self.data_processor.set_anomaly_threshold(DataType.FLOW, 0, 1000, 3.0)
    
    def connect_device(self, device: IoTDevice) -> bool:
        """
        连接IoT设备
        
        Args:
            device: IoT设备
        
        Returns:
            连接是否成功
        """
        try:
            # 模拟设备连接
            device.status = ConnectionStatus.ONLINE
            device.last_seen = datetime.now()
            
            self.connected_devices[device.device_id] = device
            self.node_info.connected_devices.append(device.device_id)
            
            return True
        except Exception as e:
            device.status = ConnectionStatus.ERROR
            return False
    
    def disconnect_device(self, device_id: str) -> bool:
        """
        断开设备连接
        
        Args:
            device_id: 设备ID
        
        Returns:
            断开是否成功
        """
        if device_id in self.connected_devices:
            self.connected_devices[device_id].status = ConnectionStatus.OFFLINE
            del self.connected_devices[device_id]
            
            if device_id in self.node_info.connected_devices:
                self.node_info.connected_devices.remove(device_id)
            
            return True
        return False
    
    def collect_sensor_data(self, device_id: str, data_type: DataType, 
                          value: float, unit: str) -> SensorData:
        """
        采集传感器数据
        
        Args:
            device_id: 设备ID
            data_type: 数据类型
            value: 数值
            unit: 单位
        
        Returns:
            传感器数据
        """
        sensor_data = SensorData(
            device_id=device_id,
            timestamp=datetime.now(),
            data_type=data_type,
            value=value,
            unit=unit,
            quality=random.uniform(0.8, 1.0),  # 模拟数据质量
            location=self.connected_devices.get(device_id, {}).location if device_id in self.connected_devices else None
        )
        
        # 处理数据
        processing_result = self.data_processor.process_sensor_data(sensor_data)
        
        # 存储到本地
        self.local_storage.append({
            'sensor_data': sensor_data,
            'processing_result': processing_result,
            'stored_at': datetime.now()
        })
        
        # 添加到同步队列
        self.sync_queue.put({
            'type': 'sensor_data',
            'data': sensor_data,
            'processing_result': processing_result
        })
        
        return sensor_data
    
    def start_processing(self):
        """
        启动数据处理
        """
        self.is_running = True
        self.processing_thread = threading.Thread(target=self._processing_loop)
        self.processing_thread.start()
    
    def stop_processing(self):
        """
        停止数据处理
        """
        self.is_running = False
        if self.processing_thread:
            self.processing_thread.join()
    
    def _processing_loop(self):
        """
        数据处理循环
        """
        while self.is_running:
            try:
                # 模拟从连接的设备采集数据
                for device_id, device in self.connected_devices.items():
                    if device.status == ConnectionStatus.ONLINE:
                        for data_type in device.data_types:
                            # 生成模拟数据
                            value = self._generate_mock_data(data_type)
                            unit = self._get_unit_for_data_type(data_type)
                            
                            self.collect_sensor_data(device_id, data_type, value, unit)
                
                # 更新节点状态
                self._update_node_status()
                
                time.sleep(1)  # 1秒采集间隔
                
            except Exception as e:
                print(f"处理循环错误: {e}")
                time.sleep(5)
    
    def _generate_mock_data(self, data_type: DataType) -> float:
        """
        生成模拟数据
        
        Args:
            data_type: 数据类型
        
        Returns:
            模拟数值
        """
        base_values = {
            DataType.TEMPERATURE: 25.0,
            DataType.PRESSURE: 10.0,
            DataType.VIBRATION: 1.0,
            DataType.FLOW: 100.0,
            DataType.VOLTAGE: 220.0,
            DataType.CURRENT: 5.0,
            DataType.SPEED: 1500.0,
            DataType.POSITION: 0.0
        }
        
        base_value = base_values.get(data_type, 0.0)
        
        # 添加随机波动
        noise = random.uniform(-0.1, 0.1) * base_value
        
        # 偶尔添加异常值
        if random.random() < 0.05:  # 5%概率异常
            noise += random.uniform(-0.5, 0.5) * base_value
        
        return max(0, base_value + noise)
    
    def _get_unit_for_data_type(self, data_type: DataType) -> str:
        """
        获取数据类型对应的单位
        
        Args:
            data_type: 数据类型
        
        Returns:
            单位字符串
        """
        units = {
            DataType.TEMPERATURE: "°C",
            DataType.PRESSURE: "bar",
            DataType.VIBRATION: "mm/s",
            DataType.FLOW: "L/min",
            DataType.VOLTAGE: "V",
            DataType.CURRENT: "A",
            DataType.SPEED: "rpm",
            DataType.POSITION: "mm"
        }
        return units.get(data_type, "")
    
    def _update_node_status(self):
        """
        更新节点状态
        """
        # 模拟资源使用情况
        self.node_info.cpu_usage = random.uniform(20, 80)
        self.node_info.memory_usage = random.uniform(30, 70)
        self.node_info.storage_usage = len(self.local_storage) / 10000 * 100  # 假设最大存储10000条记录
        
        # 更新设备状态
        for device in self.connected_devices.values():
            device.last_seen = datetime.now()
    
    def get_node_status(self) -> Dict[str, Any]:
        """
        获取节点状态
        
        Returns:
            节点状态信息
        """
        processing_stats = self.data_processor.get_processing_statistics()
        
        return {
            'node_info': asdict(self.node_info),
            'connected_devices_count': len(self.connected_devices),
            'local_storage_count': len(self.local_storage),
            'sync_queue_size': self.sync_queue.qsize(),
            'processing_statistics': processing_stats,
            'recent_alerts': self.data_processor.alerts[-10:],  # 最近10个告警
            'is_running': self.is_running
        }
    
    def sync_to_cloud(self) -> Dict[str, Any]:
        """
        同步数据到云端
        
        Returns:
            同步结果
        """
        sync_data = []
        sync_count = 0
        
        # 从同步队列获取数据
        while not self.sync_queue.empty() and sync_count < 100:  # 每次最多同步100条
            try:
                data = self.sync_queue.get_nowait()
                sync_data.append(data)
                sync_count += 1
            except Empty:
                break
        
        if sync_data:
            # 模拟云端同步
            sync_result = {
                'success': True,
                'synced_count': len(sync_data),
                'sync_time': datetime.now(),
                'node_id': self.node_info.node_id
            }
        else:
            sync_result = {
                'success': True,
                'synced_count': 0,
                'message': '无数据需要同步'
            }
        
        return sync_result

# 使用示例
print("工业物联网与边缘计算演示")
print("=" * 50)

# 创建边缘节点
edge_node_info = EdgeNode(
    node_id="EDGE001",
    node_name="生产线边缘节点1",
    location="车间A",
    ip_address="192.168.1.100",
    cpu_cores=4,
    memory_gb=8.0,
    storage_gb=256.0
)

edge_node = EdgeComputingNode(edge_node_info)

print(f"\n创建边缘节点: {edge_node_info.node_name}")
print(f"节点ID: {edge_node_info.node_id}")
print(f"位置: {edge_node_info.location}")
print(f"配置: {edge_node_info.cpu_cores}核CPU, {edge_node_info.memory_gb}GB内存")

# 创建IoT设备
devices = [
    IoTDevice(
        device_id="TEMP001",
        device_name="温度传感器1",
        device_type=DeviceType.SENSOR,
        location="生产线A",
        ip_address="192.168.1.101",
        port=502,
        protocol="Modbus",
        data_types=[DataType.TEMPERATURE],
        sampling_rate=1.0
    ),
    IoTDevice(
        device_id="PRESS001",
        device_name="压力传感器1",
        device_type=DeviceType.SENSOR,
        location="生产线A",
        ip_address="192.168.1.102",
        port=502,
        protocol="Modbus",
        data_types=[DataType.PRESSURE],
        sampling_rate=2.0
    ),
    IoTDevice(
        device_id="VIB001",
        device_name="振动传感器1",
        device_type=DeviceType.SENSOR,
        location="设备A",
        ip_address="192.168.1.103",
        port=1883,
        protocol="MQTT",
        data_types=[DataType.VIBRATION],
        sampling_rate=10.0
    )
]

print("\n连接IoT设备:")
for device in devices:
    success = edge_node.connect_device(device)
    print(f"  {device.device_name} ({device.device_id}): {'连接成功' if success else '连接失败'}")
    print(f"    类型: {device.device_type.value}, 协议: {device.protocol}")
    print(f"    数据类型: {[dt.value for dt in device.data_types]}")
    print(f"    采样频率: {device.sampling_rate}Hz")

# 启动数据处理
print("\n启动边缘数据处理...")
edge_node.start_processing()

# 运行一段时间收集数据
print("\n数据采集中...")
time.sleep(5)

# 获取节点状态
node_status = edge_node.get_node_status()
print("\n边缘节点状态:")
print(f"  连接设备数: {node_status['connected_devices_count']}")
print(f"  本地存储: {node_status['local_storage_count']}条记录")
print(f"  待同步数据: {node_status['sync_queue_size']}条")
print(f"  CPU使用率: {node_status['node_info']['cpu_usage']:.1f}%")
print(f"  内存使用率: {node_status['node_info']['memory_usage']:.1f}%")
print(f"  存储使用率: {node_status['node_info']['storage_usage']:.1f}%")

# 数据处理统计
processing_stats = node_status['processing_statistics']
if 'total_processed' in processing_stats:
    print("\n数据处理统计:")
    print(f"  总处理数据: {processing_stats['total_processed']}条")
    print(f"  异常检测: {processing_stats['anomaly_count']}条 ({processing_stats['anomaly_rate']:.1%})")
    print(f"  平均质量: {processing_stats['average_quality']:.2f}")
    print(f"  活跃规则: {processing_stats['active_rules']}个")
    print(f"  总告警数: {processing_stats['total_alerts']}个")
    
    if processing_stats['type_statistics']:
        print("\n  按类型统计:")
        for data_type, stats in processing_stats['type_statistics'].items():
            print(f"    {data_type}: {stats['count']}条, 异常{stats['anomalies']}条, 平均值{stats['avg_value']:.2f}")

# 显示最近告警
if node_status['recent_alerts']:
    print("\n最近告警:")
    for alert in node_status['recent_alerts'][-5:]:
        print(f"  [{alert['severity']}] {alert['device_id']}: {alert['message']}")
        print(f"    时间: {alert['timestamp'].strftime('%H:%M:%S')}, 值: {alert['value']:.2f}")

# 云端同步
print("\n执行云端同步...")
sync_result = edge_node.sync_to_cloud()
print(f"同步结果: {'成功' if sync_result['success'] else '失败'}")
print(f"同步数据: {sync_result['synced_count']}条")
if 'sync_time' in sync_result:
    print(f"同步时间: {sync_result['sync_time'].strftime('%Y-%m-%d %H:%M:%S')}")

# 停止处理
print("\n停止数据处理...")
edge_node.stop_processing()

print("\n工业物联网演示完成!")

29.3 边缘计算与实时处理

边缘计算概述

边缘计算是一种分布式计算架构,将计算、存储和网络服务从集中式数据中心扩展到网络边缘,更接近数据源和用户。在工业物联网中,边缘计算能够实现低延迟的实时数据处理和智能决策。

边缘计算的核心优势:

  1. 低延迟响应

    • 本地处理:数据在边缘节点本地处理,减少网络传输延迟
    • 实时决策:毫秒级响应时间,满足工业控制要求
    • 离线运行:网络中断时仍能正常工作
    • 带宽优化:减少云端数据传输,节省带宽成本
  2. 数据安全

    • 本地存储:敏感数据在边缘节点本地处理和存储
    • 数据脱敏:只上传处理后的结果,保护原始数据
    • 访问控制:边缘节点级别的安全控制
    • 合规性:满足数据本地化要求
  3. 可靠性

    • 分布式架构:单点故障不影响整体系统
    • 自主运行:边缘节点独立运行能力
    • 故障恢复:快速故障检测和恢复
    • 负载均衡:多个边缘节点分担计算负载
  4. 智能化

    • 边缘AI:在边缘节点部署机器学习模型
    • 自适应:根据环境变化自动调整参数
    • 预测分析:基于历史数据进行预测
    • 协同计算:边缘节点间的协同处理

边缘计算架构

边缘计算架构
设备层
边缘层
雾计算层
云计算层
传感器设备
执行器设备
控制器设备
智能终端
边缘网关
边缘服务器
边缘AI芯片
本地存储
区域数据中心
边缘云
CDN节点
5G基站
公有云
私有云
混合云
多云
class EdgeAIModel:
    """
    边缘AI模型
    
    功能:
    - 模型加载和管理
    - 实时推理
    - 模型更新
    - 性能监控
    """
    
    def __init__(self, model_name: str, model_type: str):
        """
        初始化边缘AI模型
        
        Args:
            model_name: 模型名称
            model_type: 模型类型
        """
        self.model_name = model_name
        self.model_type = model_type
        self.model = None
        self.is_loaded = False
        self.inference_count = 0
        self.total_inference_time = 0.0
        self.last_update = None
        self.performance_metrics = []
    
    def load_model(self, model_path: str = None) -> bool:
        """
        加载模型
        
        Args:
            model_path: 模型文件路径
        
        Returns:
            加载是否成功
        """
        try:
            # 模拟模型加载
            if self.model_type == "anomaly_detection":
                # 简单的异常检测模型(基于统计方法)
                self.model = {
                    'type': 'statistical',
                    'thresholds': {
                        'temperature': {'mean': 25.0, 'std': 5.0, 'factor': 3.0},
                        'pressure': {'mean': 10.0, 'std': 2.0, 'factor': 2.5},
                        'vibration': {'mean': 1.0, 'std': 0.5, 'factor': 2.0}
                    }
                }
            elif self.model_type == "predictive_maintenance":
                # 预测性维护模型
                self.model = {
                    'type': 'regression',
                    'features': ['temperature', 'pressure', 'vibration', 'runtime'],
                    'weights': [0.3, 0.2, 0.4, 0.1],
                    'threshold': 0.7
                }
            elif self.model_type == "quality_prediction":
                # 质量预测模型
                self.model = {
                    'type': 'classification',
                    'features': ['temperature', 'pressure', 'speed'],
                    'classes': ['good', 'defective'],
                    'decision_boundary': 0.5
                }
            
            self.is_loaded = True
            self.last_update = datetime.now()
            return True
            
        except Exception as e:
            print(f"模型加载失败: {e}")
            return False
    
    def predict(self, input_data: Dict[str, float]) -> Dict[str, Any]:
        """
        执行推理
        
        Args:
            input_data: 输入数据
        
        Returns:
            推理结果
        """
        if not self.is_loaded:
            return {'error': '模型未加载'}
        
        start_time = time.time()
        
        try:
            if self.model_type == "anomaly_detection":
                result = self._anomaly_detection(input_data)
            elif self.model_type == "predictive_maintenance":
                result = self._predictive_maintenance(input_data)
            elif self.model_type == "quality_prediction":
                result = self._quality_prediction(input_data)
            else:
                result = {'error': '未知模型类型'}
            
            # 记录性能指标
            inference_time = time.time() - start_time
            self.inference_count += 1
            self.total_inference_time += inference_time
            
            self.performance_metrics.append({
                'timestamp': datetime.now(),
                'inference_time': inference_time,
                'input_size': len(input_data),
                'success': 'error' not in result
            })
            
            # 保留最近1000条记录
            if len(self.performance_metrics) > 1000:
                self.performance_metrics = self.performance_metrics[-1000:]
            
            result['inference_time'] = inference_time
            result['model_name'] = self.model_name
            
            return result
            
        except Exception as e:
            return {'error': f'推理失败: {e}'}
    
    def _anomaly_detection(self, input_data: Dict[str, float]) -> Dict[str, Any]:
        """
        异常检测推理
        
        Args:
            input_data: 输入数据
        
        Returns:
            异常检测结果
        """
        anomalies = []
        anomaly_scores = {}
        
        for feature, value in input_data.items():
            if feature in self.model['thresholds']:
                threshold = self.model['thresholds'][feature]
                mean = threshold['mean']
                std = threshold['std']
                factor = threshold['factor']
                
                # 计算异常分数
                z_score = abs(value - mean) / std
                anomaly_scores[feature] = z_score
                
                if z_score > factor:
                    anomalies.append({
                        'feature': feature,
                        'value': value,
                        'expected_range': [mean - factor * std, mean + factor * std],
                        'anomaly_score': z_score,
                        'severity': 'HIGH' if z_score > factor * 1.5 else 'MEDIUM'
                    })
        
        return {
            'prediction': 'anomaly' if anomalies else 'normal',
            'anomalies': anomalies,
            'anomaly_scores': anomaly_scores,
            'confidence': max(anomaly_scores.values()) if anomaly_scores else 0.0
        }
    
    def _predictive_maintenance(self, input_data: Dict[str, float]) -> Dict[str, Any]:
        """
        预测性维护推理
        
        Args:
            input_data: 输入数据
        
        Returns:
            维护预测结果
        """
        features = self.model['features']
        weights = self.model['weights']
        threshold = self.model['threshold']
        
        # 计算维护需求分数
        score = 0.0
        feature_contributions = {}
        
        for i, feature in enumerate(features):
            if feature in input_data:
                # 归一化特征值
                if feature == 'temperature':
                    normalized_value = min(1.0, max(0.0, (input_data[feature] - 20) / 80))
                elif feature == 'pressure':
                    normalized_value = min(1.0, max(0.0, input_data[feature] / 50))
                elif feature == 'vibration':
                    normalized_value = min(1.0, max(0.0, input_data[feature] / 10))
                elif feature == 'runtime':
                    normalized_value = min(1.0, max(0.0, input_data.get(feature, 0) / 8760))  # 年运行小时数
                else:
                    normalized_value = 0.5
                
                contribution = weights[i] * normalized_value
                feature_contributions[feature] = contribution
                score += contribution
        
        # 预测维护需求
        maintenance_needed = score > threshold
        urgency = 'HIGH' if score > threshold * 1.2 else 'MEDIUM' if score > threshold else 'LOW'
        
        # 估算剩余时间
        if maintenance_needed:
            days_remaining = max(1, int((1.0 - score) * 30))  # 最多30天
        else:
            days_remaining = int((threshold - score) * 100)  # 基于分数差估算
        
        return {
            'prediction': 'maintenance_needed' if maintenance_needed else 'normal',
            'maintenance_score': score,
            'urgency': urgency,
            'days_remaining': days_remaining,
            'feature_contributions': feature_contributions,
            'confidence': min(1.0, abs(score - threshold) * 2)
        }
    
    def _quality_prediction(self, input_data: Dict[str, float]) -> Dict[str, Any]:
        """
        质量预测推理
        
        Args:
            input_data: 输入数据
        
        Returns:
            质量预测结果
        """
        features = self.model['features']
        boundary = self.model['decision_boundary']
        
        # 简单的质量评分计算
        quality_score = 0.0
        feature_impacts = {}
        
        # 温度影响(最优范围20-30°C)
        if 'temperature' in input_data:
            temp = input_data['temperature']
            if 20 <= temp <= 30:
                temp_impact = 1.0
            else:
                temp_impact = max(0.0, 1.0 - abs(temp - 25) / 25)
            feature_impacts['temperature'] = temp_impact
            quality_score += temp_impact * 0.4
        
        # 压力影响(最优范围8-12bar)
        if 'pressure' in input_data:
            pressure = input_data['pressure']
            if 8 <= pressure <= 12:
                pressure_impact = 1.0
            else:
                pressure_impact = max(0.0, 1.0 - abs(pressure - 10) / 10)
            feature_impacts['pressure'] = pressure_impact
            quality_score += pressure_impact * 0.3
        
        # 速度影响(最优范围1400-1600rpm)
        if 'speed' in input_data:
            speed = input_data['speed']
            if 1400 <= speed <= 1600:
                speed_impact = 1.0
            else:
                speed_impact = max(0.0, 1.0 - abs(speed - 1500) / 500)
            feature_impacts['speed'] = speed_impact
            quality_score += speed_impact * 0.3
        
        # 预测质量等级
        predicted_class = 'good' if quality_score > boundary else 'defective'
        confidence = abs(quality_score - boundary) * 2
        
        return {
            'prediction': predicted_class,
            'quality_score': quality_score,
            'confidence': min(1.0, confidence),
            'feature_impacts': feature_impacts,
            'recommendation': self._get_quality_recommendation(quality_score, feature_impacts)
        }
    
    def _get_quality_recommendation(self, quality_score: float, feature_impacts: Dict[str, float]) -> str:
        """
        获取质量改进建议
        
        Args:
            quality_score: 质量分数
            feature_impacts: 特征影响
        
        Returns:
            改进建议
        """
        if quality_score > 0.8:
            return "质量良好,继续保持当前参数"
        
        recommendations = []
        
        for feature, impact in feature_impacts.items():
            if impact < 0.7:
                if feature == 'temperature':
                    recommendations.append("调整温度至20-30°C范围")
                elif feature == 'pressure':
                    recommendations.append("调整压力至8-12bar范围")
                elif feature == 'speed':
                    recommendations.append("调整速度至1400-1600rpm范围")
        
        return "; ".join(recommendations) if recommendations else "检查设备参数设置"
    
    def get_performance_metrics(self) -> Dict[str, Any]:
        """
        获取模型性能指标
        
        Returns:
            性能指标
        """
        if not self.performance_metrics:
            return {'message': '暂无性能数据'}
        
        recent_metrics = self.performance_metrics[-100:]  # 最近100次推理
        
        inference_times = [m['inference_time'] for m in recent_metrics]
        success_rate = sum(1 for m in recent_metrics if m['success']) / len(recent_metrics)
        
        return {
            'model_name': self.model_name,
            'model_type': self.model_type,
            'total_inferences': self.inference_count,
            'average_inference_time': np.mean(inference_times),
            'max_inference_time': np.max(inference_times),
            'min_inference_time': np.min(inference_times),
            'success_rate': success_rate,
            'last_update': self.last_update,
            'is_loaded': self.is_loaded
        }

class RealTimeProcessor:
    """
    实时数据处理器
    
    功能:
    - 流式数据处理
    - 实时分析
    - 事件检测
    - 响应执行
    """
    
    def __init__(self, buffer_size: int = 1000):
        """
        初始化实时处理器
        
        Args:
            buffer_size: 缓冲区大小
        """
        self.buffer_size = buffer_size
        self.data_buffer = []
        self.ai_models = {}
        self.processing_rules = []
        self.event_handlers = {}
        self.is_processing = False
        self.processing_thread = None
        self.input_queue = Queue()
        self.output_queue = Queue()
        
    def add_ai_model(self, model: EdgeAIModel) -> bool:
        """
        添加AI模型
        
        Args:
            model: 边缘AI模型
        
        Returns:
            添加是否成功
        """
        if model.is_loaded:
            self.ai_models[model.model_name] = model
            return True
        return False
    
    def add_processing_rule(self, rule_name: str, condition: str, 
                          actions: List[str], parameters: Dict[str, Any] = None):
        """
        添加处理规则
        
        Args:
            rule_name: 规则名称
            condition: 触发条件
            actions: 执行动作列表
            parameters: 规则参数
        """
        rule = {
            'name': rule_name,
            'condition': condition,
            'actions': actions,
            'parameters': parameters or {},
            'enabled': True,
            'execution_count': 0,
            'last_executed': None
        }
        self.processing_rules.append(rule)
    
    def register_event_handler(self, event_type: str, handler_func):
        """
        注册事件处理器
        
        Args:
            event_type: 事件类型
            handler_func: 处理函数
        """
        self.event_handlers[event_type] = handler_func
    
    def start_processing(self):
        """
        启动实时处理
        """
        self.is_processing = True
        self.processing_thread = threading.Thread(target=self._processing_loop)
        self.processing_thread.start()
    
    def stop_processing(self):
        """
        停止实时处理
        """
        self.is_processing = False
        if self.processing_thread:
            self.processing_thread.join()
    
    def process_data(self, sensor_data: SensorData) -> Dict[str, Any]:
        """
        处理传感器数据
        
        Args:
            sensor_data: 传感器数据
        
        Returns:
            处理结果
        """
        # 添加到输入队列
        self.input_queue.put(sensor_data)
        
        # 如果不是实时处理模式,直接处理
        if not self.is_processing:
            return self._process_single_data(sensor_data)
        
        return {'status': 'queued', 'timestamp': datetime.now()}
    
    def _processing_loop(self):
        """
        处理循环
        """
        while self.is_processing:
            try:
                # 从输入队列获取数据
                sensor_data = self.input_queue.get(timeout=1)
                
                # 处理数据
                result = self._process_single_data(sensor_data)
                
                # 将结果放入输出队列
                self.output_queue.put(result)
                
            except Empty:
                continue
            except Exception as e:
                print(f"处理循环错误: {e}")
    
    def _process_single_data(self, sensor_data: SensorData) -> Dict[str, Any]:
        """
        处理单条数据
        
        Args:
            sensor_data: 传感器数据
        
        Returns:
            处理结果
        """
        result = {
            'sensor_data': sensor_data,
            'timestamp': datetime.now(),
            'ai_predictions': {},
            'rule_results': [],
            'events': [],
            'actions_taken': []
        }
        
        # 添加到缓冲区
        self.data_buffer.append(sensor_data)
        if len(self.data_buffer) > self.buffer_size:
            self.data_buffer.pop(0)
        
        # AI模型推理
        input_data = {
            sensor_data.data_type.value: sensor_data.value
        }
        
        # 添加历史数据特征
        recent_data = [d for d in self.data_buffer[-10:] 
                      if d.device_id == sensor_data.device_id]
        if len(recent_data) > 1:
            values = [d.value for d in recent_data]
            input_data['avg_value'] = np.mean(values)
            input_data['std_value'] = np.std(values)
            input_data['trend'] = values[-1] - values[0] if len(values) > 1 else 0
        
        # 执行AI推理
        for model_name, model in self.ai_models.items():
            try:
                prediction = model.predict(input_data)
                result['ai_predictions'][model_name] = prediction
            except Exception as e:
                result['ai_predictions'][model_name] = {'error': str(e)}
        
        # 应用处理规则
        for rule in self.processing_rules:
            if not rule['enabled']:
                continue
            
            try:
                condition_met = self._evaluate_rule_condition(rule, sensor_data, result)
                
                if condition_met:
                    rule_result = self._execute_rule_actions(rule, sensor_data, result)
                    result['rule_results'].append(rule_result)
                    
                    rule['execution_count'] += 1
                    rule['last_executed'] = datetime.now()
                    
            except Exception as e:
                result['rule_results'].append({
                    'rule_name': rule['name'],
                    'error': str(e)
                })
        
        return result
    
    def _evaluate_rule_condition(self, rule: Dict, sensor_data: SensorData, 
                               processing_result: Dict) -> bool:
        """
        评估规则条件
        
        Args:
            rule: 规则
            sensor_data: 传感器数据
            processing_result: 处理结果
        
        Returns:
            条件是否满足
        """
        condition = rule['condition']
        
        # 简化的条件评估
        if 'anomaly_detected' in condition:
            for model_name, prediction in processing_result['ai_predictions'].items():
                if prediction.get('prediction') == 'anomaly':
                    return True
        
        if 'maintenance_needed' in condition:
            for model_name, prediction in processing_result['ai_predictions'].items():
                if prediction.get('prediction') == 'maintenance_needed':
                    return True
        
        if 'quality_defective' in condition:
            for model_name, prediction in processing_result['ai_predictions'].items():
                if prediction.get('prediction') == 'defective':
                    return True
        
        # 数值条件
        if 'value >' in condition:
            threshold = float(condition.split('value >')[1].strip())
            return sensor_data.value > threshold
        
        if 'value <' in condition:
            threshold = float(condition.split('value <')[1].strip())
            return sensor_data.value < threshold
        
        return False
    
    def _execute_rule_actions(self, rule: Dict, sensor_data: SensorData, 
                            processing_result: Dict) -> Dict[str, Any]:
        """
        执行规则动作
        
        Args:
            rule: 规则
            sensor_data: 传感器数据
            processing_result: 处理结果
        
        Returns:
            执行结果
        """
        rule_result = {
            'rule_name': rule['name'],
            'actions_executed': [],
            'events_generated': [],
            'execution_time': datetime.now()
        }
        
        for action in rule['actions']:
            try:
                if action == 'send_alert':
                    event = {
                        'type': 'ALERT',
                        'device_id': sensor_data.device_id,
                        'message': rule['parameters'].get('alert_message', '设备异常'),
                        'severity': rule['parameters'].get('severity', 'MEDIUM'),
                        'timestamp': datetime.now(),
                        'data_value': sensor_data.value
                    }
                    processing_result['events'].append(event)
                    rule_result['events_generated'].append(event)
                    
                    # 调用事件处理器
                    if 'ALERT' in self.event_handlers:
                        self.event_handlers['ALERT'](event)
                
                elif action == 'trigger_maintenance':
                    event = {
                        'type': 'MAINTENANCE',
                        'device_id': sensor_data.device_id,
                        'message': '设备需要维护',
                        'urgency': rule['parameters'].get('urgency', 'MEDIUM'),
                        'timestamp': datetime.now()
                    }
                    processing_result['events'].append(event)
                    rule_result['events_generated'].append(event)
                    
                    if 'MAINTENANCE' in self.event_handlers:
                        self.event_handlers['MAINTENANCE'](event)
                
                elif action == 'adjust_parameters':
                    # 模拟参数调整
                    adjustment = {
                        'device_id': sensor_data.device_id,
                        'parameter': rule['parameters'].get('parameter', 'unknown'),
                        'new_value': rule['parameters'].get('new_value', 0),
                        'timestamp': datetime.now()
                    }
                    processing_result['actions_taken'].append(adjustment)
                
                elif action == 'log_event':
                    log_entry = {
                        'device_id': sensor_data.device_id,
                        'event_type': 'RULE_TRIGGERED',
                        'rule_name': rule['name'],
                        'data_value': sensor_data.value,
                        'timestamp': datetime.now()
                    }
                    processing_result['actions_taken'].append(log_entry)
                
                rule_result['actions_executed'].append(action)
                
            except Exception as e:
                rule_result['actions_executed'].append(f"{action} (错误: {e})")
        
        return rule_result
    
    def get_processing_status(self) -> Dict[str, Any]:
        """
        获取处理状态
        
        Returns:
            处理状态信息
        """
        return {
            'is_processing': self.is_processing,
            'buffer_size': len(self.data_buffer),
            'max_buffer_size': self.buffer_size,
            'input_queue_size': self.input_queue.qsize(),
            'output_queue_size': self.output_queue.qsize(),
            'ai_models_count': len(self.ai_models),
            'processing_rules_count': len(self.processing_rules),
            'event_handlers_count': len(self.event_handlers),
            'ai_models': {name: model.get_performance_metrics() 
                         for name, model in self.ai_models.items()}
        }

# 使用示例
print("\n边缘计算与实时处理演示")
print("=" * 50)

# 创建AI模型
print("\n创建边缘AI模型...")
anomalies_model = EdgeAIModel("anomaly_detector", "anomaly_detection")
maintenance_model = EdgeAIModel("maintenance_predictor", "predictive_maintenance")
quality_model = EdgeAIModel("quality_predictor", "quality_prediction")

# 加载模型
models = [anomalies_model, maintenance_model, quality_model]
for model in models:
    success = model.load_model()
    print(f"  {model.model_name}: {'加载成功' if success else '加载失败'}")

# 创建实时处理器
processor = RealTimeProcessor(buffer_size=500)

# 添加AI模型
for model in models:
    if model.is_loaded:
        processor.add_ai_model(model)
        print(f"  添加模型: {model.model_name}")

# 添加处理规则
print("\n配置处理规则...")
processor.add_processing_rule(
    rule_name="异常检测告警",
    condition="anomaly_detected",
    actions=["send_alert", "log_event"],
    parameters={'alert_message': '检测到设备异常', 'severity': 'HIGH'}
)

processor.add_processing_rule(
    rule_name="维护需求触发",
    condition="maintenance_needed",
    actions=["trigger_maintenance", "send_alert"],
    parameters={'urgency': 'MEDIUM', 'alert_message': '设备需要维护'}
)

processor.add_processing_rule(
    rule_name="质量异常处理",
    condition="quality_defective",
    actions=["send_alert", "adjust_parameters"],
    parameters={
        'alert_message': '产品质量异常',
        'severity': 'HIGH',
        'parameter': 'temperature',
        'new_value': 25
    }
)

processor.add_processing_rule(
    rule_name="高温告警",
    condition="value > 80",
    actions=["send_alert"],
    parameters={'alert_message': '温度过高', 'severity': 'CRITICAL'}
)

# 注册事件处理器
def alert_handler(event):
    print(f"  🚨 告警: [{event['severity']}] {event['device_id']} - {event['message']}")

def maintenance_handler(event):
    print(f"  🔧 维护: [{event['urgency']}] {event['device_id']} - {event['message']}")

processor.register_event_handler('ALERT', alert_handler)
processor.register_event_handler('MAINTENANCE', maintenance_handler)

print(f"  配置了 {len(processor.processing_rules)} 个处理规则")
print(f"  注册了 {len(processor.event_handlers)} 个事件处理器")

# 启动实时处理
print("\n启动实时处理...")
processor.start_processing()

# 模拟实时数据处理
print("\n模拟实时数据处理...")
test_data = [
    # 正常数据
    SensorData("TEMP001", datetime.now(), DataType.TEMPERATURE, 25.5, "°C"),
    SensorData("PRESS001", datetime.now(), DataType.PRESSURE, 10.2, "bar"),
    SensorData("VIB001", datetime.now(), DataType.VIBRATION, 1.1, "mm/s"),
    
    # 异常数据
    SensorData("TEMP001", datetime.now(), DataType.TEMPERATURE, 85.0, "°C"),  # 高温
    SensorData("PRESS001", datetime.now(), DataType.PRESSURE, 0.5, "bar"),    # 低压
    SensorData("VIB001", datetime.now(), DataType.VIBRATION, 8.5, "mm/s"),    # 高振动
    
    # 质量相关数据
    SensorData("TEMP001", datetime.now(), DataType.TEMPERATURE, 35.0, "°C"),  # 温度偏高
    SensorData("PRESS001", datetime.now(), DataType.PRESSURE, 15.0, "bar"),   # 压力偏高
]

processing_results = []
for i, data in enumerate(test_data):
    print(f"\n处理数据 {i+1}: {data.device_id} {data.data_type.value} = {data.value} {data.unit}")
    
    result = processor.process_data(data)
    processing_results.append(result)
    
    time.sleep(0.5)  # 模拟实时间隔

# 等待处理完成
time.sleep(2)

# 获取处理状态
status = processor.get_processing_status()
print("\n实时处理状态:")
print(f"  处理状态: {'运行中' if status['is_processing'] else '已停止'}")
print(f"  缓冲区: {status['buffer_size']}/{status['max_buffer_size']}")
print(f"  输入队列: {status['input_queue_size']}")
print(f"  输出队列: {status['output_queue_size']}")
print(f"  AI模型数: {status['ai_models_count']}")
print(f"  处理规则数: {status['processing_rules_count']}")

# AI模型性能
print("\nAI模型性能:")
for model_name, metrics in status['ai_models'].items():
    if 'total_inferences' in metrics:
        print(f"  {model_name}:")
        print(f"    推理次数: {metrics['total_inferences']}")
        print(f"    平均耗时: {metrics['average_inference_time']:.4f}秒")
        print(f"    成功率: {metrics['success_rate']:.1%}")

# 获取输出结果
print("\n处理结果摘要:")
output_results = []
while not processor.output_queue.empty():
    try:
        result = processor.output_queue.get_nowait()
        output_results.append(result)
    except Empty:
        break

if output_results:
    total_events = sum(len(r['events']) for r in output_results)
    total_actions = sum(len(r['actions_taken']) for r in output_results)
    total_rules_triggered = sum(len(r['rule_results']) for r in output_results)
    
    print(f"  处理记录数: {len(output_results)}")
    print(f"  触发事件数: {total_events}")
    print(f"  执行动作数: {total_actions}")
    print(f"  规则触发数: {total_rules_triggered}")
    
    # 显示AI预测结果统计
    prediction_stats = {}
    for result in output_results:
        for model_name, prediction in result['ai_predictions'].items():
            if model_name not in prediction_stats:
                prediction_stats[model_name] = {}
            
            pred_type = prediction.get('prediction', 'unknown')
            if pred_type not in prediction_stats[model_name]:
                prediction_stats[model_name][pred_type] = 0
            prediction_stats[model_name][pred_type] += 1
    
    print("\nAI预测统计:")
    for model_name, stats in prediction_stats.items():
        print(f"  {model_name}: {dict(stats)}")

# 停止处理
print("\n停止实时处理...")
processor.stop_processing()

print("\n边缘计算演示完成!")

## 章节总结

### 核心知识点

1. **工业物联网基础**
   - IoT架构:感知层、网络层、边缘层、平台层、应用层
   - 设备连接:传感器、执行器、控制器、网关设备
   - 通信协议:Modbus、OPC-UA、MQTT、CoAP等工业协议
   - 数据类型:温度、压力、振动、流量等工业参数

2. **边缘计算架构**
   - 分布式计算:边缘节点、雾计算、云计算的协同
   - 本地处理:实时数据处理、智能决策、离线运行
   - 资源管理:CPU、内存、存储的优化配置
   - 安全防护:边缘安全、数据加密、访问控制

3. **实时数据处理**
   - 流式处理:数据采集、预处理、实时分析
   - 异常检测:统计方法、机器学习、规则引擎
   - 事件驱动:规则配置、事件触发、动作执行
   - 性能监控:处理延迟、吞吐量、成功率

### 技术要点

1. **边缘AI部署**
   - 模型优化:模型压缩、量化、剪枝技术
   - 推理加速:专用芯片、并行计算、缓存优化
   - 模型管理:版本控制、热更新、性能监控
   - 多模型协同:异常检测、预测维护、质量预测

2. **数据处理策略**
   - 缓冲机制:数据缓存、队列管理、流量控制
   - 质量保证:数据验证、清洗、标准化
   - 实时性保证:低延迟处理、优先级调度
   - 容错机制:异常处理、故障恢复、降级策略

3. **系统集成方案**
   - 设备接入:多协议支持、设备发现、状态管理
   - 云边协同:数据同步、计算卸载、资源调度
   - 运维管理:远程监控、配置管理、故障诊断
   - 扩展性设计:水平扩展、负载均衡、弹性伸缩

### 应用前景

1. **智能制造升级**
   - 数字化工厂:全面感知、智能控制、优化决策
   - 柔性生产:快速响应、个性化定制、敏捷制造
   - 质量管控:实时检测、预防控制、持续改进
   - 设备管理:预测维护、远程诊断、生命周期管理

2. **技术发展趋势**
   - 5G+边缘计算:超低延迟、大连接、高可靠通信
   - AI芯片普及:边缘AI能力提升、成本降低
   - 数字孪生:虚实融合、仿真优化、预测分析
   - 区块链集成:数据可信、供应链透明、价值共享

3. **行业应用扩展**
   - 能源电力:智能电网、新能源管理、节能优化
   - 交通运输:智能物流、自动驾驶、交通优化
   - 环境监测:污染监控、气象预报、生态保护
   - 城市管理:智慧城市、公共安全、基础设施管理

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐