DeHighway 是一套面向高速公路运营场景的轻量级管理平台,覆盖收费站过车记录、车牌识别结果存储、异常抬杆告警以及大屏车流统计。Node.js 实现该系统的优势在于事件循环能够同时处理大量低速 I/O 请求,例如摄像头抓拍图片上传、称重设备数据写入和收费终端状态上报,而不会像传统多线程模型那样为每个连接单独分配线程导致内存快速上涨。系统采用 Express 作为 REST API 框架、Socket.IO 维护大屏长连接、Redis Streams 承担异步缓冲、MySQL 负责最终落盘,整体结构清晰,适合小团队快速迭代。

系统整体架构与数据模型
接入层由收费站工控机、车牌识别摄像头和称重设备组成,统一通过 HTTP POST 上报到 Node.js 接入服务。服务层拆成 api-server、stream-worker 和 socket-server 三个逻辑模块,可以放在同一个进程,也可以按流量拆分。数据层用 Redis Streams 做削峰,MySQL 存通行记录和收费流水。
这样拆分的主要原因是过车数据存在明显的波峰:早晚高峰可能几十个收费站同时上报,单个 Node.js 进程如果直接同步写 MySQL,连接数很快会被占满。先把消息写入 Redis Streams 能让接口在毫秒级返回,收费终端的抬杆动作不受影响,后续 worker 再按自己的节奏消费入库。
通行记录表 passage_record 的设计如下。lane_id 表示车道编号,plate_no 为车牌号,weight_kg 为称重重量,fee_cent 为收费金额,event_type 区分正常过车、异常抬杆、重复计费等情况。created_at 为事件发生时间,方便按小时聚合。
CREATE TABLE passage_record ( id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY, lane_id VARCHAR(32) NOT NULL, plate_no VARCHAR(16) NOT NULL, weight_kg DECIMAL(8,2) DEFAULT 0, fee_cent INT DEFAULT 0, event_type TINYINT NOT NULL DEFAULT 1, created_at DATETIME NOT NULL, KEY idx_lane_time (lane_id, created_at), KEY idx_plate (plate_no) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
索引方面,lane_id 和 created_at 的联合索引可以支撑按车道按时间范围查询,plate_no 单独建索引用于车辆轨迹回溯。如果后续需要统计车型分布,可以在称重设备上报时增加车型字段,并建立对应索引。
核心模块:过车采集与事件驱动处理
过车采集接口需要尽可能轻,接收到数据后只做基础校验,随后写入 Redis Streams 并立即返回。收费终端不会因为后端数据库抖动而阻塞,抬杆动作也能保持在毫秒级。下面是一个最小化的 Express 路由实现。
const express = require('express');
const { createClient } = require('redis');
const app = express();
app.use(express.json());
const redis = createClient();
redis.connect().catch(console.error);
app.post('/api/passage', async function(req, res) {
const { laneId, plateNo, weightKg, feeCent, eventType } = req.body;
if (!laneId || !plateNo) {
return res.status(400).json({ code: 400, message: '参数不完整' });
}
const record = {
laneId: laneId,
plateNo: plateNo,
weightKg: Number(weightKg || 0),
feeCent: Number(feeCent || 0),
eventType: Number(eventType || 1),
createdAt: new Date().toISOString()
};
await redis.xAdd('passage-stream', '*', { record: JSON.stringify(record) });
res.json({ code: 200, message: '已接收' });
});
app.listen(3000, function() {
console.log('api-server listening on 3000');
});这里的 Redis xAdd 并不是真正发送到收费站设备,而是写入一条流消息。流里的 record 字段使用 JSON 字符串保存,这样消费者可以自由扩展字段而不需要修改表结构。例如后续增加 ETC 状态、车型或视频片段地址时,只需要在传入 JSON 中追加。
消费端使用 Redis Streams 的消费者组机制,避免多个 worker 重复读取同一条消息。每个 worker 读取消息后解析 record,再写入 MySQL。只有写库成功后才执行 xAck 确认,如果中途异常退出,消息会进入 pending 列表等待重试。
const { createClient } = require('redis');
const mysql = require('mysql2/promise');
const redis = createClient();
redis.connect().catch(console.error);
async function startWorker() {
try {
await redis.xGroupCreate('passage-stream', 'passage-group', '$', { MKSTREAM: true });
} catch (e) {
// 消费者组可能已存在
}
while (true) {
const results = await redis.xReadGroup(
'passage-group',
'worker-' + process.pid,
{ key: 'passage-stream', id: '>' },
{ COUNT: 20, BLOCK: 1000 }
);
if (!results) continue;
for (const stream of results) {
for (const message of stream.messages) {
const payload = JSON.parse(message.message.record);
await saveToMySQL(payload);
await redis.xAck('passage-stream', 'passage-group', message.id);
}
}
}
}
async function saveToMySQL(record) {
const conn = await mysql.createConnection({
host: '127.0.0.1',
user: 'root',
password: 'password',
database: 'dehighway'
});
await conn.execute(
'INSERT INTO passage_record (lane_id, plate_no, weight_kg, fee_cent, event_type, created_at) VALUES (?, ?, ?, ?, ?, ?)',
[record.laneId, record.plateNo, record.weightKg, record.feeCent, record.eventType, new Date(record.createdAt)]
);
await conn.end();
}
startWorker().catch(console.error);异常抬杆和重复计费不能只写库,还需要实时通知值班人员。可以在 saveToMySQL 之后判断 event_type,如果属于异常类型,就 publish 到 Redis 的 alarm-channel,由 socket-server 统一接收并推送到大屏。这样采集、落库、告警三者完全解耦,后续增加短信或钉钉通知也不需要改采集接口。
事件驱动处理的关键在于每个环节只做一件事。采集接口只负责接收,stream-worker 只负责持久化,socket-server 只负责推送。Node.js 的异步模型让这些模块可以共享同一个运行环境,同时保持较低的内存占用。
实时监控与WebSocket推送
大屏需要看到最近几分钟的车流量、实时抬杆状态和异常列表。Socket.IO 适合这种场景,服务端按收费站或车道划分房间,客户端只订阅关心的车道。下面这段代码展示如何把 Redis 告警消息广播到对应车道房间。
const { createServer } = require('http');
const { Server } = require('socket.io');
const express = require('express');
const { createClient } = require('redis');
const app = express();
const httpServer = createServer(app);
const io = new Server(httpServer, { cors: { origin: '*' } });
const redis = createClient();
redis.connect().catch(console.error);
io.on('connection', function(socket) {
socket.on('subscribe-lane', function(laneId) {
socket.join('lane-' + laneId);
});
socket.on('unsubscribe-lane', function(laneId) {
socket.leave('lane-' + laneId);
});
});
async function startAlarmListener() {
const subscriber = redis.duplicate();
await subscriber.connect();
await subscriber.subscribe('alarm-channel', function(message) {
const alarm = JSON.parse(message);
io.to('lane-' + alarm.laneId).emit('alarm', alarm);
});
}
httpServer.listen(3001);
startAlarmListener().catch(console.error);客户端连接后发送 subscribe-lane 事件,传入车道编号即可加入房间。当某个车道产生异常抬杆时,服务端只给该房间内的连接发送 alarm 事件,大屏前端根据 laneId 更新对应模块的颜色和告警文本。如果将所有告警无差别广播给所有客户端,几十块大屏会收到大量无关消息,前端还需要做过滤,网络和 CPU 都不划算。
在实际部署中,socket-server 需要独立扩展以支撑更多长连接。可以接入 Redis 的适配器,多个 socket-server 进程共享房间信息和订阅关系,这样客户端无论连接到哪个进程,都能收到正确的广播。Redis 适配器内部通过发布订阅模拟跨进程房间,开发阶段可以直接使用单进程,上线时再增加配置。
除了告警推送,还可以定期发送统计快照。服务端每隔 5 秒从 Redis 读取最近一小时的分钟级计数,组装成数组后 emit 给大屏。快照推送比前端轮询 HTTP 接口更省资源,也更容易控制刷新频率。
数据统计与性能优化
按小时、按天统计车流量不能只靠实时查 MySQL,因为 passage_record 表每天百万级增长,大屏每 5 秒查一次聚合 SQL 会产生大量慢查询。更合适的做法是用 Redis 计数器加定时任务。每次过车事件在写入 Stream 后对相应时间片做 INCR,过期时间设为 7 天,这样热数据始终在内存里。
计数器的键可以设计成 traffic:hour:20250318:15 和 traffic:day:20250318。小时键适用于大屏近实时展示,天键用于日报。增加计数的代码只需两行,可以放在采集接口里,也可以放在 stream-worker 消费成功后,以避免未落盘数据被提前统计。
const statsKey = 'traffic:hour:' + record.createdAt.slice(0, 13).replace(':', '') + ':' + new Date(record.createdAt).getUTCHours().toString().padStart(2, '0');
await redis.incr(statsKey);
await redis.expire(statsKey, 604800);这段代码中的时间处理是为了演示,实际项目应根据服务器时区统一转换为北京时间,否则凌晨边界会偏。推荐使用 dayjs 的时区插件处理,避免 Native Date 的 UTC 和本地时区混用问题。
定时任务用 node-cron 每小时把 Redis 中的小时计数写入统计表,写完后保留原始计数,方便大屏继续读取当天数据。如果担心进程重启导致任务不执行,可以引入分布式锁,只让一个 worker 执行聚合写入。Node.js 的单线程特性让计数器累加不用考虑锁,但多进程部署时 Redis 的原子 INCR 依然能保证数据正确。
const cron = require('node-cron');
const mysql = require('mysql2/promise');
cron.schedule('0 * * * *', async function() {
const keys = await redis.keys('traffic:hour:*');
for (const key of keys) {
const value = await redis.get(key);
const parts = key.split(':');
const hourKey = parts[1] + parts[2] + parts[3];
await saveStats(hourKey, Number(value));
}
});除了 Redis 计数器,还可以在 MySQL 侧建立汇总表 census_hour,避免每次统计都扫描明细表。census_hour 只保留时间片、车道编号和过车数量三个维度,查询效率远高于对 passage_record 做 GROUP BY。这种方式属于典型的读写分离思路:明细表负责记录,汇总表负责分析。
Node.js 在性能优化上还要注意 CPU 密集计算不要阻塞事件循环。如果需要对车牌识别结果做模糊匹配或图像处理,应放到 worker_threads 或独立服务中执行。高速公路管理系统的核心流量路径是写入和推送,保持主线程轻量才能让 Redis 命令、Socket.IO 事件和 HTTP 响应保持低延迟。
Node.js高速公路管理系统WebSocket推送修改时间:2026-09-23 17:58:07