物联网 · 盛世宏博 · 系统解耦架构设计

档案库房"八防"(防火、防盗、防潮、防光、防尘、防高温、防低温、防虫)一体化系统在实际落地中,普遍存在将传感器采集、逻辑判断、设备驱动全部耦合在单一控制器或单一软件进程中的做法。初期看似简洁,但随着系统规模扩大、设备种类增多、联动逻辑复杂化,耦合架构的隐性代价逐渐暴露:
耦合架构的典型问题 | 具体表现 |
|---|---|
故障扩散 | 一个传感器通信异常导致整个控制循环阻塞,空调、除湿机全部失联 |
升级困难 | 新增一种传感器协议需要重新编译整个系统,停机升级影响全天候运行 |
调试复杂 | 环境数据异常时无法快速判断是感知层问题还是控制层问题 |
扩展性差 | 不同库房的设备品牌、协议各异,每接入一个新设备就要改主控逻辑 |
责任边界模糊 | 感知不准导致控制误动作时,难以界定是传感器问题还是控制算法问题 |
解耦的核心目标:将"感知"与"控制"拆分为两个独立运行的模块,通过标准化接口通信。感知模块只负责"看",控制模块只负责"动",两者之间的决策逻辑作为独立中间层,各司其职、互不阻塞。
┌─────────────────────────────────────────────────────────────────────┐
│ 解耦三层架构 │
│ │
│ ┌─────────────────────┐ 标准数据接口 ┌────────────────────┐│
│ │ 环境感知模块 │ ───────────────▶ │ 决策引擎层 ││
│ │ │ │ ││
│ │ · 传感器协议适配 │ 控制指令接口 │ · 八防规则引擎 ││
│ │ · 数据采集与校准 │ ◀─────────────── │ · 联动逻辑判断 ││
│ │ · 异常检测与标记 │ │ · 防冲突仲裁 ││
│ │ · 本地缓存 │ │ · 安全边界守护 ││
│ └─────────────────────┘ └─────────┬──────────┘│
│ │ │
│ ┌─────────────────────┐ 设备驱动接口 │ │
│ │ 设备控制模块 │ ◀──────────────────────────┘ │
│ │ │ │
│ │ · 设备协议适配 │ 状态反馈接口 │
│ │ · 指令队列与重试 │ ─────────────────────▶ 决策引擎层 │
│ │ · 互锁保护 │ │
│ │ · 手动/自动切换 │ │
│ └─────────────────────┘ │
│ │
│ 关键设计原则: │
│ · 感知模块与控制模块之间无直接调用关系,全部通过决策引擎中转 │
│ · 任一模块崩溃不影响其他模块独立运行(感知断联时控制维持最后状态) │
│ · 模块间通信采用消息队列,支持异步、缓冲、重放 │
└─────────────────────────────────────────────────────────────────────┘# sensing_module.py
from dataclasses import dataclass, field
from typing import Dict, List, Optional, Callable
from enum import Enum
import time
import json
class SensorStatus(Enum):
NORMAL = "normal"
WARNING = "warning" # 读数可疑但未确认故障
FAULT = "fault" # 确认故障
CALIBRATING = "calibrating"
@dataclass
class SensorReading:
sensor_id: str
sensor_type: str # temperature / humidity / smoke / water / pm25 / voc / lux / door
zone_id: str
value: float
unit: str
timestamp: float
status: SensorStatus = SensorStatus.NORMAL
confidence: float = 1.0 # 置信度 0~1
raw_data: Optional[dict] = None
@dataclass
class SensorConfig:
sensor_id: str
sensor_type: str
zone_id: str
protocol: str # modbus / bacnet / lorawan / tcp / analog
address: str
poll_interval: float = 30.0
calibration_offset: float = 0.0
valid_range: tuple = (0, 100)
fault_threshold: int = 3 # 连续异常次数触发故障
class SensingModule:
"""环境感知模块——只负责采集、校准、标记,不做任何控制决策"""
def __init__(self, message_bus=None):
self.sensors: Dict[str, SensorConfig] = {}
self.latest_readings: Dict[str, SensorReading] = {}
self.consecutive_errors: Dict[str, int] = {}
self.message_bus = message_bus
self._running = False
def register_sensor(self, config: SensorConfig):
"""注册传感器"""
self.sensors[config.sensor_id] = config
self.consecutive_errors[config.sensor_id] = 0
def poll_once(self, sensor_id: str) -> Optional[SensorReading]:
"""单次轮询单个传感器"""
config = self.sensors.get(sensor_id)
if not config:
return None
try:
raw_value = self._read_from_device(config)
value = raw_value + config.calibration_offset
# 量程校验
if value < config.valid_range[0] or value > config.valid_range[1]:
self.consecutive_errors[sensor_id] += 1
status = SensorStatus.WARNING
confidence = 0.3
else:
self.consecutive_errors[sensor_id] = 0
status = SensorStatus.NORMAL
confidence = 1.0
# 故障判定
if self.consecutive_errors[sensor_id] >= config.fault_threshold:
status = SensorStatus.FAULT
confidence = 0.0
reading = SensorReading(
sensor_id=sensor_id,
sensor_type=config.sensor_type,
zone_id=config.zone_id,
value=value,
unit=self._get_unit(config.sensor_type),
timestamp=time.time(),
status=status,
confidence=confidence,
raw_data={"raw": raw_value, "protocol": config.protocol}
)
self.latest_readings[sensor_id] = reading
# 发布到消息总线(供决策引擎消费)
if self.message_bus:
self.message_bus.publish("sensing/reading", reading)
return reading
except Exception as e:
self.consecutive_errors[sensor_id] += 1
if self.consecutive_errors[sensor_id] >= config.fault_threshold:
reading = SensorReading(
sensor_id=sensor_id,
sensor_type=config.sensor_type,
zone_id=config.zone_id,
value=0,
unit=self._get_unit(config.sensor_type),
timestamp=time.time(),
status=SensorStatus.FAULT,
confidence=0.0
)
self.latest_readings[sensor_id] = reading
if self.message_bus:
self.message_bus.publish("sensing/fault", reading)
return None
def poll_all(self) -> Dict[str, SensorReading]:
"""轮询所有传感器"""
results = {}
for sensor_id in self.sensors:
reading = self.poll_once(sensor_id)
if reading:
results[sensor_id] = reading
return results
def get_readings_by_zone(self, zone_id: str) -> Dict[str, SensorReading]:
"""按区域获取最新读数"""
return {
sid: r for sid, r in self.latest_readings.items()
if r.zone_id == zone_id and r.status != SensorStatus.FAULT
}
def get_readings_by_type(self, sensor_type: str) -> Dict[str, SensorReading]:
"""按类型获取最新读数"""
return {
sid: r for sid, r in self.latest_readings.items()
if r.sensor_type == sensor_type and r.status != SensorStatus.FAULT
}
def _read_from_device(self, config: SensorConfig) -> float:
"""根据协议读取设备数据(实际项目中对接具体协议库)"""
# 此处为示意,实际对接 modbus/lorawan/bacnet 等
if config.protocol == "modbus":
return self._modbus_read(config.address)
elif config.protocol == "bacnet":
return self._bacnet_read(config.address)
elif config.protocol == "tcp":
return self._tcp_read(config.address)
else:
return 0.0
def _modbus_read(self, address: str) -> float:
# 实际实现:modbus_tk / pymodbus
return 22.5 # placeholder
def _bacnet_read(self, address: str) -> float:
# 实际实现:bacpypes
return 50.0 # placeholder
def _tcp_read(self, address: str) -> float:
# 实际实现:socket 通讯
return 22.0 # placeholder
def _get_unit(self, sensor_type: str) -> str:
units = {
"temperature": "°C", "humidity": "%RH", "smoke": "ppm",
"water": "bool", "pm25": "μg/m³", "voc": "mg/m³",
"lux": "lux", "door": "bool", "static": "V", "co2": "ppm"
}
return units.get(sensor_type, "unknown")# control_module.py
from dataclasses import dataclass, field
from typing import Dict, List, Optional, Callable
from enum import Enum
import time
import queue
class DeviceStatus(Enum):
IDLE = "idle"
RUNNING = "running"
ERROR = "error"
OFFLINE = "offline"
LOCKED = "locked" # 被互锁保护
class DeviceMode(Enum):
AUTO = "auto" # 自动模式(接收决策引擎指令)
MANUAL = "manual" # 手动模式(现场优先)
MAINTENANCE = "maintenance" # 维护模式(禁用自动指令)
@dataclass
class DeviceCommand:
device_id: str
command: str # start / stop / set_param
params: dict = field(default_factory=dict)
priority: int = 5 # 1最高
timestamp: float = 0
source: str = "decision_engine" # 指令来源
@dataclass
class DeviceState:
device_id: str
device_type: str # ac / dehumidifier / humidifier / purifier / exhaust / fresh_air / siren / door_lock
zone_id: str
status: DeviceStatus
current_params: dict = field(default_factory=dict)
last_command: Optional[DeviceCommand] = None
last_response: str = ""
mode: DeviceMode = DeviceMode.AUTO
interlocks: List[str] = field(default_factory=list) # 互锁设备列表
class ControlModule:
"""设备控制模块——只负责接收指令、驱动设备、反馈状态"""
def __init__(self, message_bus=None):
self.devices: Dict[str, DeviceState] = {}
self.command_queue: queue.PriorityQueue = queue.PriorityQueue()
self.message_bus = message_bus
self._running = False
self.interlock_table: Dict[str, List[str]] = {
"dehumidifier": ["humidifier"],
"humidifier": ["dehumidifier"],
"ac_cool": ["heater"],
"heater": ["ac_cool"],
"fresh_air": ["exhaust"],
}
def register_device(self, device_id: str, device_type: str,
zone_id: str, interlocks: List[str] = None):
"""注册设备"""
self.devices[device_id] = DeviceState(
device_id=device_id,
device_type=device_type,
zone_id=zone_id,
status=DeviceStatus.IDLE,
interlocks=interlocks or []
)
def submit_command(self, cmd: DeviceCommand) -> bool:
"""提交控制指令(来自决策引擎)"""
device = self.devices.get(cmd.device_id)
if not device:
return False
# 手动/维护模式拦截
if device.mode == DeviceMode.MAINTENANCE:
return False
if device.mode == DeviceMode.MANUAL and cmd.source != "local_manual":
return False
# 互锁检查
if not self._check_interlock(device, cmd):
device.status = DeviceStatus.LOCKED
return False
cmd.timestamp = cmd.timestamp or time.time()
self.command_queue.put((cmd.priority, cmd))
return True
def execute_loop(self):
"""指令执行循环(独立线程运行)"""
while self._running:
try:
priority, cmd = self.command_queue.get(timeout=1.0)
self._execute_command(cmd)
except queue.Empty:
continue
def _execute_command(self, cmd: DeviceCommand):
"""执行单条指令"""
device = self.devices.get(cmd.device_id)
if not device:
return
try:
# 根据设备类型调用对应驱动
result = self._drive_device(device, cmd)
if result:
device.status = DeviceStatus.RUNNING if cmd.command == "start" else DeviceStatus.IDLE
device.last_command = cmd
device.last_response = "ok"
# 反馈状态到消息总线
if self.message_bus:
self.message_bus.publish("control/state_change", {
"device_id": device.device_id,
"status": device.status.value,
"params": device.current_params
})
else:
device.status = DeviceStatus.ERROR
device.last_response = "execution_failed"
except Exception as e:
device.status = DeviceStatus.ERROR
device.last_response = str(e)
def _drive_device(self, device: DeviceState, cmd: DeviceCommand) -> bool:
"""设备驱动层(实际项目中对接具体协议)"""
if device.device_type == "ac":
return self._drive_ac(device, cmd)
elif device.device_type == "dehumidifier":
return self._drive_dehumidifier(device, cmd)
elif device.device_type == "humidifier":
return self._drive_humidifier(device, cmd)
elif device.device_type == "purifier":
return self._drive_purifier(device, cmd)
elif device.device_type == "fresh_air":
return self._drive_fresh_air(device, cmd)
elif device.device_type == "exhaust":
return self._drive_exhaust(device, cmd)
elif device.device_type == "siren":
return self._drive_siren(device, cmd)
else:
return False
def _drive_ac(self, device: DeviceState, cmd: DeviceCommand) -> bool:
# 对接空调控制器
if cmd.command == "start":
device.current_params["power"] = "on"
if "mode" in cmd.params:
device.current_params["mode"] = cmd.params["mode"]
elif cmd.command == "stop":
device.current_params["power"] = "off"
elif cmd.command == "set_param":
device.current_params.update(cmd.params)
return True
def _drive_dehumidifier(self, device: DeviceState, cmd: DeviceCommand) -> bool:
if cmd.command == "start":
device.current_params["power"] = "on"
elif cmd.command == "stop":
device.current_params["power"] = "off"
return True
def _drive_humidifier(self, device: DeviceState, cmd: DeviceCommand) -> bool:
if cmd.command == "start":
device.current_params["power"] = "on"
elif cmd.command == "stop":
device.current_params["power"] = "off"
return True
def _drive_purifier(self, device: DeviceState, cmd: DeviceCommand) -> bool:
if cmd.command == "start":
device.current_params["power"] = "on"
if "level" in cmd.params:
device.current_params["level"] = cmd.params["level"]
elif cmd.command == "stop":
device.current_params["power"] = "off"
return True
def _drive_fresh_air(self, device: DeviceState, cmd: DeviceCommand) -> bool:
if cmd.command == "start":
device.current_params["power"] = "on"
if "speed" in cmd.params:
device.current_params["speed"] = cmd.params["speed"]
elif cmd.command == "stop":
device.current_params["power"] = "off"
return True
def _drive_exhaust(self, device: DeviceState, cmd: DeviceCommand) -> bool:
if cmd.command == "start":
device.current_params["power"] = "on"
elif cmd.command == "stop":
device.current_params["power"] = "off"
return True
def _drive_siren(self, device: DeviceState, cmd: DeviceCommand) -> bool:
if cmd.command == "start":
device.current_params["power"] = "on"
elif cmd.command == "stop":
device.current_params["power"] = "off"
return True
def _check_interlock(self, device: DeviceState, cmd: DeviceCommand) -> bool:
"""互锁检查"""
all_interlocks = set(device.interlocks)
for dev_id, dev_state in self.devices.items():
if dev_id != device.device_id:
for il in dev_state.interlocks:
if il == device.device_id:
all_interlocks.add(dev_id)
if cmd.command == "start":
for other_id in all_interlocks:
other = self.devices.get(other_id)
if other and other.status == DeviceStatus.RUNNING:
return False
return True
def set_device_mode(self, device_id: str, mode: DeviceMode):
"""设置设备模式(手动/自动/维护)"""
device = self.devices.get(device_id)
if device:
device.mode = mode
def get_device_states(self, zone_id: str = None) -> Dict[str, DeviceState]:
"""获取设备状态"""
if zone_id:
return {k: v for k, v in self.devices.items() if v.zone_id == zone_id}
return dict(self.devices)# decision_engine.py
from dataclasses import dataclass
from typing import Dict, List, Optional
import time
@dataclass
class Rule:
rule_id: str
trigger_type: str # sensor_type
trigger_condition: dict # {"operator": ">", "value": 24}
target_devices: List[dict] # [{"device_id": "ac_01", "command": "start", "params": {...}}]
priority: int = 5
cooldown: float = 300 # 冷却时间(秒)
enabled: bool = True
class DecisionEngine:
"""决策引擎——解耦感知与控制的中间层"""
def __init__(self, sensing_module, control_module, message_bus=None):
self.sensing = sensing_module
self.control = control_module
self.message_bus = message_bus
self.rules: Dict[str, Rule] = {}
self.last_trigger_time: Dict[str, float] = {}
self._running = False
def register_rule(self, rule: Rule):
self.rules[rule.rule_id] = rule
def evaluate(self, reading) -> List[dict]:
"""评估单条感知数据,生成控制指令"""
if reading.status.value == "fault":
return [] # 故障传感器不参与决策
triggered_actions = []
for rule in self.rules.values():
if not rule.enabled:
continue
if rule.trigger_type != reading.sensor_type:
continue
# 条件判断
condition = rule.trigger_condition
value = reading.value
matched = False
if condition.get("operator") == ">" and value > condition["value"]:
matched = True
elif condition.get("operator") == "<" and value < condition["value"]:
matched = True
elif condition.get("operator") == ">=" and value >= condition["value"]:
matched = True
elif condition.get("operator") == "<=" and value <= condition["value"]:
matched = True
elif condition.get("operator") == "==" and value == condition["value"]:
matched = True
if matched:
# 冷却检查
now = time.time()
last = self.last_trigger_time.get(rule.rule_id, 0)
if now - last < rule.cooldown:
continue
for target in rule.target_devices:
cmd = DeviceCommand(
device_id=target["device_id"],
command=target["command"],
params=target.get("params", {}),
priority=rule.priority,
source="decision_engine"
)
success = self.control.submit_command(cmd)
if success:
triggered_actions.append({
"rule_id": rule.rule_id,
"device_id": target["device_id"],
"command": target["command"]
})
self.last_trigger_time[rule.rule_id] = now
return triggered_actions
def run_loop(self):
"""主循环:消费感知消息,驱动决策"""
while self._running:
# 从消息总线获取最新读数
reading = self.message_bus.consume("sensing/reading")
if reading:
actions = self.evaluate(reading)
if actions and self.message_bus:
self.message_bus.publish("decision/actions", actions)
time.sleep(0.1)def load_eight_protection_rules(engine: DecisionEngine):
"""加载八防联动规则"""
# 防高温
engine.register_rule(Rule(
rule_id="anti_high_temp",
trigger_type="temperature",
trigger_condition={"operator": ">", "value": 24},
target_devices=[{"device_id": "ac_01", "command": "start",
"params": {"mode": "cool", "setpoint": 22}}],
priority=3, cooldown=300
))
# 防低温
engine.register_rule(Rule(
rule_id="anti_low_temp",
trigger_type="temperature",
trigger_condition={"operator": "<", "value": 14},
target_devices=[{"device_id": "heater_01", "command": "start"}],
priority=3, cooldown=300
))
# 防潮
engine.register_rule(Rule(
rule_id="anti_high_humi",
trigger_type="humidity",
trigger_condition={"operator": ">", "value": 60},
target_devices=[{"device_id": "dehumidifier_01", "command": "start"}],
priority=3, cooldown=600
))
# 防火
engine.register_rule(Rule(
rule_id="anti_fire",
trigger_type="smoke",
trigger_condition={"operator": ">", "value": 0.5},
target_devices=[
{"device_id": "ac_01", "command": "stop"},
{"device_id": "fresh_air_01", "command": "stop"},
{"device_id": "exhaust_01", "command": "start"},
{"device_id": "siren_01", "command": "start"}
],
priority=1, cooldown=60
))
# 防水
engine.register_rule(Rule(
rule_id="anti_water",
trigger_type="water",
trigger_condition={"operator": "==", "value": 1},
target_devices=[
{"device_id": "humidifier_01", "command": "stop"},
{"device_id": "dehumidifier_01", "command": "stop"},
{"device_id", "alarm_panel_01", "command": "start"}
],
priority=1, cooldown=60
))
# 防光
engine.register_rule(Rule(
rule_id="anti_light",
trigger_type="lux",
trigger_condition={"operator": ">", "value": 300},
target_devices=[{"device_id": "curtain_01", "command": "start",
"params": {"position": "closed"}}],
priority=4, cooldown=120
))
# 防尘
engine.register_rule(Rule(
rule_id="anti_dust",
trigger_type="pm25",
trigger_condition={"operator": ">", "value": 75},
target_devices=[
{"device_id": "purifier_01", "command": "start", "params": {"level": "high"}},
{"device_id": "fresh_air_01", "command": "stop"}
],
priority=3, cooldown=600
))
# 防盗
engine.register_rule(Rule(
rule_id="anti_theft",
trigger_type="door",
trigger_condition={"operator": "==", "value": 1},
target_devices=[
{"device_id": "camera_01", "command": "start", "params": {"rec": True}},
{"device_id": "siren_01", "command": "start"}
],
priority=1, cooldown=30
))# message_bus.py
from typing import Dict, List, Callable
import queue
import json
class MessageBus:
"""轻量级消息总线——模块间异步通信"""
def __init__(self):
self._topics: Dict[str, List[Callable]] = {}
self._queues: Dict[str, queue.Queue] = {}
def subscribe(self, topic: str, callback: Callable):
if topic not in self._topics:
self._topics[topic] = []
self._topics[topic].append(callback)
def publish(self, topic: str, message):
if topic in self._topics:
for callback in self._topics[topic]:
try:
callback(message)
except Exception:
pass
# 持久化队列(供轮询消费)
if topic not in self._queues:
self._queues[topic] = queue.Queue(maxsize=1000)
try:
self._queues[topic].put_nowait(message)
except queue.Full:
pass
def consume(self, topic: str, timeout=1.0):
if topic not in self._queues:
return None
try:
return self._queues[topic].get(timeout=timeout)
except queue.Empty:
return None感知模块 消息总线 决策引擎 控制模块
│ │ │ │
│── sensor_reading ──▶│ │ │
│ │── callback ────────▶│ │
│ │ │── evaluate ──▶ rules│
│ │ │── cmd ────────────▶│
│ │ │ │── drive device
│ │ │◀── state_change ───│
│◀── (可选)校准反馈 │ │ │改造前:感知与控制耦合在单台PLC中,一次温湿度传感器RS485总线短路导致PLC死机,整个库房的空调、除湿机全部停摆4小时。
改造方案:
sensing_module.py,通过RS485轮询所有传感器。 control_module.py,通过继电器模块和Modbus TCP驱动空调、除湿机、新风。 效果:
指标 | 改造前 | 改造后 |
|---|---|---|
单点故障影响范围 | 全系统 | 仅该传感器/设备 |
新增传感器停机时间 | 2~4小时 | 0(热插拔注册) |
联动响应时间 | 5~10秒 | <2秒 |
调试定位时间 | 平均30分钟 | 平均5分钟(模块隔离) |
问题 | 排查方向 | 解决 |
|---|---|---|
感知数据不更新 | 感知模块进程是否存活、传感器通信是否正常 | 检查进程状态,用poll_once()单点测试 |
控制指令无响应 | 控制模块指令队列是否积压、设备是否在线 | 查看队列深度,检查设备连接状态 |
联动规则不触发 | 规则条件是否匹配、冷却时间是否未过 | 打印evaluate()返回值,检查last_trigger_time |
设备互锁误触发 | 互锁表配置是否正确 | 检查interlock_table和check_interlock()逻辑 |
消息总线消息丢失 | 队列是否满、消费者是否阻塞 | 增大队列容量,检查消费者处理耗时 |
解耦后延迟增大 | 网络延迟、消息序列化开销 | 改用共享内存或本地Unix socket |
关键词:档案八防,一体化系统,解耦设计,环境感知模块,设备控制模块,决策引擎,消息总线,联动控制,互锁保护,模块化架构
标签:#档案八防 #一体化系统 #解耦设计 #环境感知 #设备控制 #决策引擎 #消息总线 #联动控制 #互锁保护 #模块化架构
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。