NodeJS实现SerialPort回调转Promise:等待指定ID响应
串口SerialPort监听转Promise响应的实现方案
核心实现思路
通过维护待处理请求映射表,将每个请求ID与对应的Promise的resolve/reject函数绑定,串口监听到响应后匹配ID触发对应Promise,同时结合串行写入队列和超时机制保证可靠性:
- 用
Map存储待处理请求,键为请求ID,值为包含resolve、reject的对象 - 串口数据监听器解析响应ID,找到对应请求后触发resolve/reject并清理映射
- 为每个请求设置超时,避免内存泄漏
- 结合现有MQ串行写入逻辑,确保命令发送顺序可控,但不强制响应顺序
代码示例(Node.js环境)
串口通信封装类
const SerialPort = require('serialport'); const Readline = require('@serialport/parser-readline'); class SerialCommunication { constructor(portConfig) { this.port = new SerialPort(portConfig); this.parser = this.port.pipe(new Readline({ delimiter: '\r\n' })); // 存储待处理请求:key=请求ID,value={ resolve, reject } this.pendingRequests = new Map(); // 写入消息队列,保证串行写入 this.writeQueue = []; this.isWriting = false; // 绑定数据监听和错误处理 this.parser.on('data', (data) => this.handleResponse(data)); this.port.on('error', (err) => this.handlePortError(err)); } // 解析串口响应,触发对应Promise handleResponse(rawData) { try { // 假设响应为JSON格式,需根据实际设备协议调整解析逻辑 const response = JSON.parse(rawData.trim()); const requestId = response.id; if (!requestId) return; const pending = this.pendingRequests.get(requestId); if (!pending) return; // 根据响应是否包含错误触发不同回调 response.error ? pending.reject(new Error(response.error)) : pending.resolve(response.data); this.pendingRequests.delete(requestId); } catch (err) { console.error('响应解析失败:', err); } } // 端口错误时,终止所有待处理请求 handlePortError(err) { for (const [_, { reject }] of this.pendingRequests) { reject(new Error(`串口错误: ${err.message}`)); } this.pendingRequests.clear(); } // 串行写入串口(基于MQ实现) async writeIntoSerial(command) { return new Promise((resolve, reject) => { this.writeQueue.push({ command, resolve, reject }); if (!this.isWriting) this.processWriteQueue(); }); } // 处理写入队列 async processWriteQueue() { this.isWriting = true; while (this.writeQueue.length) { const { command, resolve, reject } = this.writeQueue.shift(); try { await this.port.write(command); await this.port.drain(); resolve(); } catch (err) { reject(err); } } this.isWriting = false; } // 等待对应ID的响应,返回Promise waitUntilCommunicationDone(requestId, timeout = 5000) { return new Promise((resolve, reject) => { // 避免重复请求同一ID if (this.pendingRequests.has(requestId)) { reject(new Error(`请求ID ${requestId}已处于待处理状态`)); return; } // 设置超时清理 const timeoutId = setTimeout(() => { this.pendingRequests.delete(requestId); reject(new Error(`请求 ${requestId} 超时(${timeout}ms)`)); }, timeout); // 封装带超时清理的回调 const wrappedResolve = (data) => { clearTimeout(timeoutId); resolve(data); }; const wrappedReject = (err) => { clearTimeout(timeoutId); reject(err); }; this.pendingRequests.set(requestId, { resolve: wrappedResolve, reject: wrappedReject }); }); } }
调用示例
// 初始化串口 const serial = new SerialCommunication({ path: '/dev/ttyUSB0', baudRate: 9600 }); // 封装带响应的命令发送函数 async function sendCommand(commandId, commandContent) { // 构造带ID的命令(需匹配设备协议) const commandStr = JSON.stringify({ id: commandId, content: commandContent }); // 串行写入命令 await serial.writeIntoSerial(commandStr); // 等待对应ID的响应 return await serial.waitUntilCommunicationDone(commandId); } // 外部任意调用,自动匹配响应 sendCommand('temp_001', 'get') .then(temp => console.log('当前温度:', temp)) .catch(err => console.error('温度请求失败:', err)); sendCommand('humi_002', 'get') .then(humi => console.log('当前湿度:', humi)) .catch(err => console.error('湿度请求失败:', err));
关键注意事项
- 协议匹配:需根据设备实际的命令/响应格式调整
handleResponse中的解析逻辑(比如非JSON格式的字符串分割、二进制解析) - 线程安全:Node.js单线程环境下,
pendingRequests的所有操作都是同步执行的,不存在多线程竞争问题,无需额外锁机制 - 内存泄漏防护:超时机制会自动清理未完成的请求,端口错误时也会清空所有待处理请求
- 请求ID唯一性:需保证每个请求的ID全局唯一,避免不同请求的响应被错误匹配
内容的提问来源于stack exchange,提问作者Gipyo.Choi
相关产品推荐
相关产品推荐

