
化工安全监测的核心矛盾在于数据量大且时序严格,而处理延迟和系统崩溃比功能缺失更危险。Node.js 单线程模型配合异步非阻塞 I/O,天生适合处理大量并发的传感器连接和消息分发。我们并不打算拿它去替代已有的 DCS 或 SIS 系统,而是在其上方架设一层可编程的安全壳,把 Modbus TCP、MQTT、OPC UA 等协议的实时数据汇聚起来,做二次判定和多渠道告警。这样一来,你既能保留既有控制系统的完整性,又能用低成本的方式获得定制化报警逻辑和移动端推送能力。
化工数据入口:多协议接入与协议转换
生产过程里常见的数据源五花八门,有支持 Modbus TCP 的质量流量计,有走 MQTT 的无线温度探头,还有只能提供 RS-485 串口的老旧设备。Node.js 生态里有一个利器——node-red 的底层库其实可以被剥离出来单独使用,但我们更倾向于直接使用 modbus-serial、mqtt 和 node-opcua 这些原生模块,这样能让逻辑更轻便可控。把不同协议的数据统一成内部 JSON 格式,是后续处理的第一步。举个例子,从 Modbus 寄存器读到的两组 16 位整数需要拼成一个 32 位浮点数,这种事情放在接入层处理,可以让报警引擎只面对干净的结构化数据。
考虑一个典型的反应车间:三台反应釜各配一个温度变送器和一个压力变送器,它们通过 Modbus TCP 上传数据;车间内还有一套可燃气体探测器走 MQTT 协议。我们可以在 Node.js 进程里为每个协议创建独立的连接管理器,每个管理器内部维护一个设备地址到标签名的映射表。当收到原始报文时,解析模块立即把十六进制数据转为物理量,并打上时间戳和来源标签,塞进一个进程内的环形缓冲区。这样做的好处是,即使下层的 DCS 系统重启或网络抖动,我们的监测层也因为内存缓冲而不至于丢失关键数据点。
// Modbus TCP 连接与解析示例
const ModbusRTU = require("modbus-serial");
const client = new ModbusRTU();
async function connectModbus(ip, port) {
await client.connectTCP(ip, { port });
client.setID(1);
console.log(`已连接到 Modbus 设备 ${ip}:${port}`);
}
async function readReactorData() {
// 读取保持寄存器,起始地址 0,长度 4(两个16位整数表示温度浮点数)
const data = await client.readHoldingRegisters(0, 4);
// 将两个16位整数拼成32位IEEE 754浮点数
const buffer = Buffer.alloc(4);
buffer.writeUInt16BE(data.data[0], 0);
buffer.writeUInt16BE(data.data[1], 2);
const temperature = buffer.readFloatBE(0);
// 同理读取压力值
buffer.writeUInt16BE(data.data[2], 0);
buffer.writeUInt16BE(data.data[3], 2);
const pressure = buffer.readFloatBE(0);
return { temperature, pressure, timestamp: Date.now() };
}
对于 MQTT 接入,可以直接利用 mqtt 库连接企业的 broker,订阅相关主题。需要注意的是,化工现场经常有多个环境,比如生产区和罐区,它们可能部署在不同的 MQTT 服务器上。我们可以抽象一个订阅管理器,用配置文件驱动,动态加载需要监听的 topic 列表,这样新增一个监测点的时候不用改动代码,只需更新配置。整个接入层一定要做好断线重连和心跳检测,Node.js 的异步特性在这里体现得淋漓尽致,一个连接断开不会阻塞其他连接的数据接收。
数据管道与环形缓冲区设计
不同来源的数据到达频率不同,温度变送器可能每秒更新一次,而气体探测器每 5 秒才发布一条消息。如果我们直接把原始数据推送给报警引擎,会导致大量的无效判断和 CPU 抖动。所以需要一个时间窗口对齐与降采样机制。Node.js 里可以用数组实现简单的环形缓冲区,但面对高频率写入时,推荐使用 circular-buffer 这个轻量 npm 包,它可以设置固定容量,并在写满后自动覆盖最旧的数据,避免内存无限增长。
我们围绕每个关键测点创建一个缓冲区实例,比如反应釜 R101 的温度是一个 60 秒的时间窗口,每秒一个点,容量就是 60。每次接收到新数据点,不仅存入缓冲区,还要计算出最近 5 秒、15 秒的变化率。变化率本身就是一个强有力的预警指标,比如温度在 15 秒内上升超过 2 摄氏度,即使还没达到硬报警阈值,也应该触发预警告。Node.js 的单线程模型让我们可以直接在数据进入入库或判断前同步地执行这些统计计算,延时可以忽略不计。
// 简易环形缓冲区与变化率计算
class SensorBuffer {
constructor(capacity) {
this.buffer = new Array(capacity);
this.capacity = capacity;
this.writeIndex = 0;
this.isFull = false;
}
push(valueObj) {
this.buffer[this.writeIndex] = valueObj;
this.writeIndex = (this.writeIndex + 1) % this.capacity;
if (this.writeIndex === 0) this.isFull = true;
}
getRecent(count) {
if (!this.isFull && this.writeIndex < count) return [];
const start = (this.writeIndex - count + this.capacity) % this.capacity;
if (start < this.writeIndex) {
return this.buffer.slice(start, this.writeIndex);
} else {
return this.buffer.slice(start).concat(this.buffer.slice(0, this.writeIndex));
}
}
rateOfChange(seconds) {
const recent = this.getRecent(seconds);
if (recent.length < 2) return 0;
const first = recent[0].value;
const last = recent[recent.length - 1].value;
return (last - first) / ((recent[recent.length - 1].timestamp - first.timestamp) / 1000);
}
}
// 使用示例:R101 温度窗口 60 秒
const r101TempBuffer = new SensorBuffer(60);
// 每次接到数据后
r101TempBuffer.push({ value: 145.3, timestamp: Date.now() });
const change5s = r101TempBuffer.rateOfChange(5);
if (change5s > 2.0) {
console.log('R101 温度 5 秒内快速上升,当前速率:', change5s, '°C/s');
}
数据管道中还应该加入一个简单的去抖动逻辑,防止某一帧因为电磁干扰出现离谱的数值。常见做法是与上一帧比较,如果差值超过正常物理变化范围的 5 倍就直接丢弃,并记录一条质量位异常日志。这部分代码很短,放在数据推入缓冲区之前执行即可。
报警引擎的实现:阈值与逻辑组合
化工安全的报警规则很少是单一阈值这么简单。举个实际例子,进料阀关闭时反应釜压力应该缓慢下降,如果此时压力反而上升,就必须立即切断进料并声光报警,哪怕压力值本身还没到高高限。这种逻辑需要把开关量信号和模拟量信号联动判断。我们可以用一个小型的规则引擎来实现,而 Node.js 里完全不需要引入像 Drools 那样的重量级框架,自己写一个解释器就足够了。把每条规则定义成一个 JSON 对象,包含条件表达式、严重级别、动作指令以及报警文本模板。条件表达式中可以引用当前测点的实时值、开关量状态和变化率,运行时用 eval() 有安全风险,改用 new Function() 并传入受控的上下文变量会更稳妥。
// 规则定义示例
const rules = [
{
id: 'R1',
desc: 'R101 温度高高报',
condition: 'getValue("R101.temp") > 180',
level: 'H',
action: 'alarm',
message: '反应釜R101温度超过180°C,请立即处理'
},
{
id: 'R2',
desc: '压力异常高报(阀门关闭时)',
condition: 'getValue("R101.pressure") > 1.2 && getValue("R101.feedValve") === 0',
level: 'HH',
action: 'emergency',
message: '进料阀已关闭但压力仍高于1.2MPa,存在泄漏或反应失控风险'
},
{
id: 'R3',
desc: '温度快速上升',
condition: 'getRate("R101.temp",15) > 1.5',
level: 'H',
action: 'preAlarm',
message: '反应釜R101温度15秒内上升超过1.5°C/s,关注反应强度'
}
];
// 规则引擎执行器
class RuleEngine {
setContext(getValueFn, getRateFn) {
this.getValue = getValueFn;
this.getRate = getRateFn;
}
evaluate(rules) {
const matched = [];
for (const rule of rules) {
try {
const fn = new Function('getValue', 'getRate', `return (${rule.condition});`);
if (fn(this.getValue, this.getRate)) {
matched.push(rule);
}
} catch (err) {
console.error(`规则 ${rule.id} 求值失败:`, err.message);
}
}
return matched;
}
}
报警引擎每收到一帧新的全局状态快照就执行一轮规则匹配。这里说的全局状态快照,是把所有测点的最新值、开关量和算法因子汇聚到一个对象里,由数据管道定时(比如每 200 毫秒)生成。匹配出的报警规则需要进入一个抑制与升级模块:同一个测点的同一级别报警在 10 秒内不重复推送,避免短信轰炸;如果从 H 级上升到 HH 级,则立刻升级并覆盖之前的屏蔽。抑制状态可以保存在内存 Map 中,Key 由测点 ID 加报警级别组成。
实时推送与告警分发:WebSocket + MQTT 双通道
化工安全监测系统的告警输出必须到达中控大屏、现场声光报警器,以及不在岗人员的企业微信或钉钉。Node.js 里最拿手的 WebSocket 可以承担与前端监控画面的双向通信,而 MQTT 则适合把报警报文发给公共广播系统或短信猫服务端。使用 ws 模块搭建 WebSocket 服务器非常简单,当报警引擎产生一条未屏蔽的报警时,我们把它序列化成统一的 JSON 告警报文,同时广播给所有已连接的客户端和一个内部的 MQTT 发布模块。
对于大屏展示,前端通常需要一张实时更新的报警列表,以及每个测点的趋势图数据。WebSocket 推送时可以采用频道设计,客户端在连接时发送订阅消息,指定需要哪些车间的数据,服务端据此过滤推送内容,减少不必要的数据传输。针对移动端推送,我们可以在 Node.js 中集成企业微信机器人 Webhook,将告警内容格式化成 markdown 并通过 HTTP POST 发送。需要注意的是,所有外部 HTTP 请求都应该设置超时和错误重试,防止因为第三方接口响应慢而阻塞整个事件循环。可以利用 axios 的拦截器统一处理全局超时 3 秒,超时后直接记录失败日志并尝试备用通道。
// WebSocket 广播与告警分发简例
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8081 });
wss.on('connection', (ws) => {
ws.on('message', (msg) => {
try {
const { type, workshop } = JSON.parse(msg);
if (type === 'subscribe') {
ws.workshop = workshop; // 记录订阅的车间
}
} catch(e) {}
});
});
function broadcastAlarm(alarmData) {
const payload = JSON.stringify(alarmData);
wss.clients.forEach((client) => {
// 如果客户端订阅了全局或指定车间,才发送
if (client.readyState === WebSocket.OPEN &&
(!client.workshop || client.workshop === alarmData.workshop)) {
client.send(payload);
}
});
// 同时发布到 MQTT 告警主题
mqttClient.publish('alarm/all', payload, { qos: 1 });
}
分发的幂等性也要考虑。某些告警可能因为抑制窗口刚过就再次触发,导致重复推送。我们在发送告警前,可以为每条告警生成一个基于规则 ID、测点值和触发时间的哈希,存入一个有过期时间的 Set。如果短时间内再次产生相同哈希,就丢弃。这种轻量的去重机制足以应对绝大部分重复抖动。
日志、快照回放与自愈机制
安全系统自身的可靠性同样重要。Node.js 进程必须能够记录所有报警事件以及触发前后一段时间的原始数据,方便事后追溯。我们使用内置的 fs 模块和日志库 winston,将报警事件以 JSON 行格式写入滚动日志文件。同时也应定期(例如每 10 分钟)将环形缓冲区里的所有数据快照存入时序数据库或者就是简单的 CSV 文件。一旦系统需要重启或发生崩溃,启动时可以读取最近的快照文件,快速恢复缓冲区状态,避免冷启动造成几十秒的监控盲区。
自愈机制集中在容错处理上。Node.js 的 process.on('uncaughtException') 和 unhandledRejection 事件必须被捕获,记录完致命日志后,最稳妥的做法是调用 process.exit(1) 并借助进程守护工具(如 PM2 或 systemd)自动拉起新进程。在退出前,要尽可能将内存中尚未持久化的报警和最新数据状态刷到磁盘。PM2 的 Cluster 模式下,我们还可以开启多个工作进程,一个进程专门负责 WebSocket 通信,另一个负责数据处理和报警逻辑,避免因某个客户端的异常断开影响核心监测。
另外,建议实现一个内部的健康检查 HTTP 端点,返回当前所有连接的协议状态、缓冲区长度和最后数据时间戳。中控室可以接入这个端点,一旦发现某个 Modbus 连接断开超过 20 秒,立即在屏幕上给出黄色提示,而不是等到报警时才发现通讯中断。这一整套机制,让 Node.js 不仅成为报警中枢,也变成了系统自诊断的一个窗口。
化工行业的特殊性要求任何技术选型都不能盲目赶时髦,但 Node.js 在非阻塞 I/O、事件驱动逻辑和快速迭代方面的优势,确实能让安全监测层变得更加轻盈和可靠。你不需要用 Java 或 C# 写一套庞然大物,几百行 JavaScript 搭配合理的架构就能撑起一个车间级的安全保护伞。