实时推送不一定都要上Redis或消息队列。在小规模单机项目中,SQLite的触发器和WAL模式配合WebSocket长连接,可以构建一套足够稳定的实时通知系统。核心思路是:业务表发生写入时,由SQLite触发器自动把变更写入一张轻量的日志表;Node.js服务端以固定间隔读取增量日志,再通过WebSocket广播给所有在线客户端。这个过程不需要引入独立中间件,部署和运维成本都很低。

一、先判断这套方案是否适合你的项目
SQLite作为嵌入式数据库,写入性能受单文件锁限制,但它足够支撑每秒几百次写入,对于内部看板、设备状态监控、实验室数据采集等场景完全够用。如果系统本身已经部署在单台服务器上,WebSocket连接数在几百个以内,采用SQLite加触发器日志的方案会比额外部署Redis更省心。
需要明确的是,这套方案依赖固定间隔轮询日志表,实时性通常在200毫秒到1秒之间。若业务要求毫秒级延迟,或者需要多实例部署、水平扩展,那么更适合Redis的发布订阅或PostgreSQL的LISTEN/NOTIFY机制。但这类需求往往伴随更多运维工作,小型项目并不划算。
另外,使用SQLite做实时推送,意味着写入必须集中在一个进程或单实例服务中。如果你已经有一个Node.js服务负责所有数据库写入,那么把推送逻辑放在同一个进程里非常自然。若写入来自多个进程或不同语言的服务,维护触发器虽然可行,但广播层需要独立出来,复杂度会上升。
二、用触发器把变更记录下来
实时推送的第一步是准确捕获数据变化。SQLite支持在表上创建AFTER INSERT、AFTER UPDATE和AFTER DELETE触发器,我们可以在触发器中向专门的变更日志表写入一条记录。这样业务代码不需要关心通知逻辑,任何插入操作都会自动登记。
下面先创建业务表sensor_data和变更日志表data_changes。日志表只保存发生变化的表名、行ID和变更类型,具体数据仍然从原表读取,避免冗余。
PRAGMA journal_mode = WAL;
CREATE TABLE IF NOT EXISTS sensor_data (
id INTEGER PRIMARY KEY AUTOINCREMENT,
device_id TEXT NOT NULL,
value REAL NOT NULL,
created_at TEXT DEFAULT (datetime('now','localtime'))
);
CREATE TABLE IF NOT EXISTS data_changes (
id INTEGER PRIMARY KEY AUTOINCREMENT,
table_name TEXT NOT NULL,
row_id INTEGER NOT NULL,
change_type TEXT NOT NULL,
created_at TEXT DEFAULT (datetime('now','localtime'))
);
CREATE TRIGGER IF NOT EXISTS sensor_data_insert_trigger
AFTER INSERT ON sensor_data
BEGIN
INSERT INTO data_changes(table_name, row_id, change_type)
VALUES ('sensor_data', NEW.id, 'INSERT');
END;
这段SQL同时开启了WAL模式。WAL允许读写并发,写入操作不会阻塞读取,这对轮询程序非常重要。如果数据库默认使用DELETE日志模式,连续写入时查询变更日志可能碰到锁等待,影响广播及时性。
触发器创建完成后,向sensor_data插入一条数据,data_changes表就会自动多出一行记录。后续服务端只关注data_changes中ID大于上次已处理位置的行,即可实现增量读取。如果还要处理更新和删除,可以按照同样方式增加AFTER UPDATE和AFTER DELETE触发器,在日志中记录对应的change_type。
三、服务端轮询与WebSocket广播实现
服务端使用Node.js配合better-sqlite3和ws两个库。better-sqlite3是同步API,执行效率高,代码也更容易阅读。轮询间隔设为500毫秒,每次查询新增日志,再根据日志中的row_id回查业务表拿到完整数据,统一以JSON格式广播。
WebSocket服务启动后会维护一个clients集合,每个新连接建立时先发送欢迎消息和历史数据,方便前端立即展示当前状态。之后每隔500毫秒检查数据变化,一条条广播给所有在线客户端。
const WebSocket = require('ws');
const Database = require('better-sqlite3');
const db = new Database('./monitor.db');
db.pragma('journal_mode = WAL');
db.pragma('busy_timeout = 5000');
const wss = new WebSocket.Server({ port: 8080 });
let lastSentId = db.prepare('SELECT COALESCE(MAX(id), 0) AS maxId FROM data_changes').get().maxId;
function broadcast(obj) {
const msg = JSON.stringify(obj);
wss.clients.forEach(client => {
if (client.readyState === WebSocket.OPEN) {
client.send(msg);
}
});
}
setInterval(() => {
const rows = db.prepare('SELECT * FROM data_changes WHERE id > ? ORDER BY id ASC').all(lastSentId);
if (rows.length > 0) {
rows.forEach(row => {
const detail = db.prepare('SELECT * FROM sensor_data WHERE id = ?').get(row.row_id);
broadcast({ type: 'data-change', table: row.table_name, changeType: row.change_type, data: detail });
});
lastSentId = rows[rows.length - 1].id;
}
}, 500);
wss.on('connection', ws => {
ws.send(JSON.stringify({ type: 'welcome', message: 'connected' }));
const latest = db.prepare('SELECT * FROM sensor_data ORDER BY id DESC LIMIT 50').all();
ws.send(JSON.stringify({ type: 'history', data: latest }));
});
这里使用lastSentId记录已经发送过的最大日志ID,避免重复推送。进程重启时lastSentId从数据库中当前最大值开始,意味着重启期间产生的历史日志不会自动补发,但初始历史查询可以兜底获取最新50条。若需要更严格的事件补发,可以把lastSentId持久化到配置表或文件中。
还需要注意,轮询函数里对每条日志都执行了一次业务表回查。如果一次批量插入产生了大量日志,可以用IN查询一次性取出所有相关业务行,减少数据库往返次数。当前示例保持简单,方便理解整体流程。
四、前端接收消息与断线重连
前端部分不依赖任何第三方库,直接用浏览器原生WebSocket API。连接成功后监听message事件,根据type字段区分实时数据和历史数据。实时数据渲染到列表顶部,历史数据可以批量展示,帮助用户快速了解已有记录。
<!DOCTYPE html>
<html>
<head>
<meta charset="UTF-8">
<title>实时监控</title>
</head>
<body>
<h2>设备数据实时推送</h2>
<ul id="messageList"></ul>
<script>
const ws = new WebSocket('ws://127.0.0.1:8080');
const list = document.getElementById('messageList');
ws.onmessage = function(event) {
const msg = JSON.parse(event.data);
if (msg.type === 'data-change') {
const li = document.createElement('li');
li.textContent = '设备:' + msg.data.device_id + ',数值:' + msg.data.value;
list.appendChild(li);
}
if (msg.type === 'history') {
msg.data.forEach(item => {
const li = document.createElement('li');
li.textContent = '历史数据:' + item.device_id + ' - ' + item.value;
list.appendChild(li);
});
}
};
ws.onclose = function() {
setTimeout(() => {
location.reload();
}, 3000);
};
</script>
</body>
</html>
断线重连逻辑放在onclose事件中,延迟3秒后刷新页面。简单场景下刷新页面已经足够,如果需要更好的体验,可以把重连逻辑改为重新创建WebSocket对象,并维护心跳检测。页面中所有标签名称在正文里都应写成转义形式,例如<script>标签,以免被浏览器解析。
生产环境中建议把WebSocket地址改成当前页面所在主机的IP或域名,而不是固定使用127.0.0.1。如果前端页面与服务端不在同一台机器上,还需要确认防火墙是否放行对应端口。
五、上线前需要注意的优化点
日志表data_changes会随时间不断增长,如果不清理最终会拖慢查询。可以增加一个定时任务,每天删除超过7天的日志行:DELETE FROM data_changes WHERE created_at < datetime('now', '-7 days')。也可以基于ID保留最近10万条,删除策略根据实际写入频率决定。
另一个容易忽视的是WebSocket心跳。代理服务器或NAT设备可能在一段时间无数据后断开空闲连接,因此可以在服务端增加每30秒发送一次ping帧的逻辑,客户端收到后自动回复pong。ws库内置了心跳检测方法,也可以在message事件中约定心跳消息。
如果业务表更新频繁,还可以为data_changes的created_at字段建立索引,加速清理和按时间查询。对于更高实时性要求,可以缩短轮询间隔到100毫秒,但要观察CPU占用和数据库读压力。一般来说,500毫秒已经能覆盖大多数监控场景。
最后,多终端同时在线时,广播使用同一个WebSocket连接集合即可。若某个客户端频繁断线或响应缓慢,可以考虑为每个连接设置发送超时,并在异常时主动关闭连接,避免无效连接占用资源。