ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Node.js流之read(size)实战:精准拆包、解决粘包半包与串口设备解析

Node.js流之read(size)实战:精准拆包、解决粘包半包与串口设备解析 1. 为什么我放弃了 data 事件read(size) 解决的三个实际问题1.1 只靠 data 事件做协议解析有多痛苦先说一段我自己的经历。之前做一个工业设备对接设备通过 TCP 持续往服务端推数据数据格式是固定的 18 字节一帧2 字节帧头、4 字节重量值、8 字节时间戳、4 字节校验位。刚开始我用最直觉的写法监听 data 事件把 chunk 一个个拼起来let buffer Buffer.alloc(0); socket.on(data, (chunk) { buffer Buffer.concat([buffer, chunk]); while (buffer.length 18) { const frame buffer.subarray(0, 18); buffer buffer.subarray(18); handleFrame(frame); } });这段代码本身没毛病但它暴露了一个非常现实的问题data 事件的 chunk 大小完全由操作系统和 TCP 层决定跟你协议里的帧边界没有任何关系。实测中我遇到的情况有三种一次 data 事件来了 62 字节里面包含 3 个半帧一次 data 事件只来了 5 字节连一帧都不够更恶心的是由于 Nagle 算法和内核缓冲区两个 data 事件之间还会出现你等了几十毫秒才凑够一帧的假象。也就是说用 data 事件做协议解析本质上是在不间断地拼接 切割你得自己维护一个残留缓冲每来一个 chunk 都要做一次 Buffer.concat 和 subarray。这个逻辑本身不难但只要协议一复杂比如帧头带长度字段、允许多个 frame 粘连、需要考虑半包和粘包data 事件的方案就会迅速变成一团乱麻。1.2 readable 模式给了你手动按需读取的能力Node.js 的可读流Readable其实有两种工作模式流动模式flowing mode和暂停模式paused mode。你一旦监听了 data 事件或者调用了 pipe()流就会进入流动模式数据像自来水一样哗哗流过来你只能被动接收没法控制一次接多少。而当你监听 readable 事件时流会进入暂停模式数据会先堆积在内部缓冲区里由你主动用 read() 去取。这个差异非常关键。以前很多教程只会教你监听 readable 事件后调用 read()read() 不传参数就是读全部缓存但很少有人强调read() 是可以传 size 参数的传了 size就意味着你要求流一次性返回指定字节数的数据。这个能力才是做协议解析时真正能救命的东西。还是用刚才那个 18 字节一帧的例子。如果改用 readable read(18)循环读帧的逻辑就变得异常清晰socket.on(readable, () { let frame; while ((frame socket.read(18)) ! null) { handleFrame(frame); } });你不需要自己维护 buffer不需要考虑这次来了几个帧、是不是半帧read(18) 会帮你从内部缓冲里精准切出 18 字节。如果内部缓冲不够 18 字节read(18) 会返回 null数据继续留在缓冲里等底层数据到了、凑够了 18 字节readable 事件会再次触发你再接着读。整个拆包的过程被 read(size) 给标准化了。1.3 三个实际问题read(size) 恰好对症结合我的实战经验read(size) 至少解决了三个 data 事件很难优雅处理的问题第一个是半包问题。数据不足一帧时data 方案需要你自己把残留数据存起来read(size) 方案里数据不足时直接返回 null内部缓冲天然帮你暂存不需要额外变量。第二个是粘包问题。数据里包含多个帧时data 方案要 while 循环切割read(size) 方案同样可以用 while 循环但循环条件里的 read(18) 是按帧消费不会出现把下一个帧的字节切给上一个帧的情况因为每次只取 18 字节取完它内部会继续尝试取下一个 18 字节。第三个是背压控制。流动模式下数据来了你就得处理如果处理不过来内存迟早被撑爆暂停模式下你不调 read() 数据就堆在缓冲区里堆到 highWaterMark 之后底层会停止向流里继续拉数据这等于给了你一个天然的刹车。配合 read(size)你可以做到每次只处理一个帧处理完再取下一个从而精细控制消费速度。2. read(size) 的底层行为与边界参数一次性讲透2.1 内部缓冲区是怎么拼接数据的很多人第一次接触 read(size) 时会误以为流内部的缓冲区就是一个大一维字节数组read(18) 就是从里面 slice 18 个字节出来。实际不是。Node.js 的 Readable 内部用的是 BufferList 结构它把多个独立的 chunk 像链表一样串起来。每个 chunk 可能是 4096 字节可能是 512 字节也可能是 3 字节。当调用 read(size) 时内部逻辑会遍历这个 BufferList跨多个 chunk 拼接出恰好 size 大小的数据再返回。举个例子内部缓冲里现在有三个 chunk第一个 10 字节第二个 5 字节第三个 20 字节。你调用 read(18)它会取出第一个 chunk 的全部 10 字节再取第二个 chunk 的 5 字节再从第三个 chunk 里取 3 字节凑成一个 18 字节的 Buffer 返回。第三个 chunk 还剩下 17 字节继续留在缓冲里。这个过程对使用者来说是透明的你完全不需要关心数据的物理存储方式。但这也解释了为什么 read(size) 在处理帧跨多个 chunk时特别顺手它内部已经帮你做了跨 chunk 的拼接你拿到手的就是一个完整的帧不需要自己在应用层再拼一次。2.2 返回 null 的条件与处理方式read(size) 返回 null在绝大多数情况下只有一个原因当前内部缓冲区里的数据量不够 size 字节。但有两个细节需要注意。第一个细节是返回 null 不代表流结束了。流可能还在源源不断地产数据只是现在这一刻缓冲区内数据不足。你应该继续监听 readable 事件等下一次触发时再尝试 read(size)。所以正确的循环写法一定是stream.on(readable, () { let chunk; while ((chunk stream.read(18)) ! null) { // 处理一个完整的帧 } });注意这个 while 循环它会在缓冲区内有足够数据时把所有 18 字节的帧都读完直到读到 null 才退出。不能只写一次 read(18)因为一次 readable 事件触发时缓冲区里可能有多个帧的数据。第二个细节是流结束时EOF也会触发一次 readable 事件。这时候如果缓冲区内还有残留数据且不足 size 字节read(size) 同样返回 null。你需要用 read() 不传参数的方式把残留数据取出来做收尾处理这个我在后面的实战案例里会详细演示。2.3 read(0)、read()、read(-1) 分别是什么行为除了常见的 read(size) 和 read()还有一些边界调用容易被忽略read(0)不会真的读取任何数据返回值一定是 null。但它的作用不是什么都不做而是会触发一次底层的 _read 调用让流去底层拉取数据。你可以用它来催促数据尽快到达常用于某些需要主动拉取的场景。注意它不会消费缓冲区里已有的数据。read()不传参数返回内部缓冲区里当前所有可读的数据。如果缓冲区为空且流已经结束返回 null。read(-1)和 read() 行为一致也是返回所有可读数据。负数参数会被当作未指定 size处理。这几个边界行为在官方文档里写得很简略实战中你会在一个地方真正用到它们当你读了若干个定长帧最后发现缓冲区里还剩一点不够一帧的数据时你用 read() 把残留数据全部取出来保存到自己的变量里等下次数据到了再拼。这就是 2.2 里说的收尾处理。2.4 objectMode 模式下的 size 是无效的还需要强调一个容易踩的坑如果流是 objectMode对象模式也就是流里传输的不是 Buffer 而是 JavaScript 对象那么 read(size) 里的 size 参数会被忽略。每次 read() 只能读出一个对象或者返回 null。字节控流只对二进制流有意义设计上它就不是为对象流准备的。我曾经看到有人试图用 objectMode 流做按字节解析死活读不到想要的数量最后排查半天才发现是 objectMode 的问题。如果你的协议解析手里拿着的是对象流要么改成二进制流要么自己维护数据拼接read(size) 帮不了你。3. 定长帧解析实战read(size) 的正确打开方式3.1 从 0 写一个定长帧解析器实战永远是理解 API 最好的方式。假设我们有一个设备通过串口或者 TCP 发送数据每帧固定 18 字节前 2 字节是帧头 0xAA 0x55中间 4 字节是重量值小端序后面是时间戳和校验。需求就是实时解析出每一帧并且把重量值打印出来。用 read(size) 实现的核心逻辑如下const { Readable } require(stream); class FrameParser extends Readable { constructor(options) { super(options); this.frameLength 18; this.remain Buffer.alloc(0); } _read() {} // 外部把数据推进来 pushData(data) { if (!this.push(data)) { // push 返回 false 表示内部缓冲满了这里可以做一些流控处理 } } _readFrameFromBuffer() { let frame this.read(this.frameLength); if (frame null) { // 缓冲区里的数据不够一帧把剩余数据全部取出保存 this.remain this.read(); return null; } return frame; } }这个类继承了 Readable实现了一个空的 _read 方法因为数据是通过外部 pushData 主动推入的不需要 Readable 自己从底层拉数据。重点是父类内部的缓冲区机制以及 read(this.frameLength) 的精准读取。3.2 在 readable 事件里循环读取帧数据有了上面的 FrameParser 类外部调用就变得非常清爽const parser new FrameParser({ highWaterMark: 1024 }); parser.on(readable, () { let frame; while ((frame parser.read(18)) ! null) { const weight frame.readInt32LE(2); console.log(当前重量: ${weight} 克); } }); // 模拟从串口收到数据每次 push 的 chunk 大小故意不一致 parser.pushData(Buffer.from([0xAA, 0x55, 0x34, 0x12])); parser.pushData(Buffer.from([0x00, 0x00, 0x78, 0x56, 0x34, 0x12, 0x00, 0x00, 0x00, 0x00, 0x01, 0x02, 0x03, 0x04])); parser.pushData(Buffer.from([0xAA, 0x55, 0x00, 0x05]));第一次 push 了 4 字节远不够 18 字节read(18) 返回 null这 4 字节被留在内部缓冲区。第二次 push 了 12 字节此时内部缓冲一共 16 字节依然不够 18 字节继续留在缓冲。第三次 push 了 4 字节内部缓冲总共 20 字节read(18) 成功取走 18 字节还剩 2 字节。如果第三次 push 的数据实际上是下一个帧的开头那么在 while 循环里继续调用 read(18) 会发现缓冲只剩 2 字节返回 null然后我们把剩余 2 字节通过 read() 保存到 remain 变量中等待下一次数据到来再拼接。3.3 一个关键细节拿到 Buffer 后要立即消费用 read(size) 拿到的 Buffer 看似安全其实背后有一个隐藏的池化内存问题。Node.js 为了性能会从预分配的 Buffer 池里切出内存给流使用。这意味着你拿到的 Buffer 可能指向的是一个更大的共享内存块。虽然 Node.js 内部对流的 read 做了处理大多数情况下你拿到的 Buffer 已经是独立切片但我在实际开发中养成了一个习惯凡是跨异步边界使用的帧数据一律先拷贝到自己的 Buffer 里。比如parser.on(readable, () { let frame; while ((frame parser.read(18)) ! null) { const copy Buffer.from(frame); // 拷贝一份防止后续被覆盖 processFrame(copy); } });尤其是当你把 frame 传给异步任务、或者把 frame 存进队列延后处理时拷贝几乎是必须的。我在这里踩过坑帧数据还没来得及处理下一次 read 就把底层的 BufferPool 内容覆盖了导致读取到的重量值忽大忽小。新手很难排查这种 bug因为它不是必现的只有在内存复用压力大的时候才会冒出来。3.4 流结束时的残留数据收尾还有一个容易被忽略的场景设备断开连接、或者文件读完了流会触发 end 事件。此时如果缓冲区里还残留着不足一帧的数据read(18) 会返回 null你永远等不到下一次 readable 事件了。所以必须在 end 事件里做收尾parser.on(end, () { const tail parser.read(); if (tail tail.length 0) { // 最后一段数据不完整说明设备可能中途断开 console.log(警告最后一帧不完整剩余 ${tail.length} 字节); } });这段代码能帮你快速定位数据传了一半连接断了这类问题。有一次我在排查嵌入式设备上报数据丢失的问题就是靠这个收尾逻辑发现设备在异常断电前只发了半帧过来。4. 精准控流的坑编码、缓冲、流结束与内存那些事4.1 setEncoding 之后read(size) 的 size 单位变了很多人在解析文本流时喜欢调用 setEncoding(utf8)这样 read() 返回的就是字符串而不是 Buffer。但如果你同时用 read(size) 做按数量读取就会踩一个大坑设置了 encoding 之后size 参数的含义变成了字符数而不是字节数。UTF-8 是一种变长编码一个中文字符占 3 字节一个英文字符占 1 字节。如果你 read(18)你以为取 18 字节实际取的是 18 个字符换算下来可能是 30 字节。这在二进制协议解析里会造成致命错乱。我给自己定的规矩是做二进制协议解析时绝对不调用 setEncoding。你不需要字符串你需要的是 Buffer是字节是精确定位到偏移量的数据。setEncoding 是给人类读文本用的不是给机器拆协议用的。4.2 highWaterMark 和 read(size) 的关系highWaterMark 是 Readable 流内部缓冲区的水位线当缓冲数据量超过这个值底层就会暂停从数据源拉取更多数据。它和 read(size) 之间有一个非常微妙的相互作用。如果 highWaterMark 设置得比较小比如 64 字节而你一帧要读 18 字节多数情况下没问题。但如果你一帧要读 1024 字节highWaterMark 却只有 256 字节那么可能发生一件反直觉的事你调用 read(1024)内部缓冲区只有 256 字节按前面说的逻辑应该返回 null。但实际上 Node.js 在处理这种情况时会临时调整内部状态触发底层 _read 继续拉数据直到缓冲区够 1024 字节或者流结束。这里的关键是不要让 read(size) 的 size 超过 highWaterMark 太多。我习惯将 highWaterMark 设置成 size 的 2~4 倍比如帧大小是 18 字节就把 highWaterMark 设为 1024帧大小是 4096就设为 16384。这样既能保证 read(size) 有足够缓冲可以拼接又不会让内存无限膨胀。参数配置参考帧大小size建议 highWaterMark说明64 字节以内1024小帧高频场景避免缓冲过大导致延迟1KB 左右4096常规协议包足够拼接4KB~16KB16384 或更高大帧场景缓冲不足会频繁触发补充读取4.3 残留缓冲的重复复制陷阱回到协议解析的场景。每次帧不够时我们用 read() 把残留数据取出来存到 remain 变量等下次数据来了再 Buffer.concat([remain, newData])。这个逻辑简单但性能上有隐患如果残留数据长期积累、每次都要 concat那么会发生大量的内存拷贝。比如极端情况下一帧 1MB设备每次只传 500KB那你每 500KB 就要做一次 concat拷贝的累计成本是 O(n^2)。数据量小的时候无所谓数据量大了会发现 CPU 和 GC 压力暴涨。针对这个问题我常用的优化手段是预分配固定大小的 Buffer用游标offset管理写入位置而不是反复 concat。比如维护一个 2 倍帧大小的缓冲数组把收到的数据直接写入对应位置等小区里数据量达到一帧再从固定偏移读取。这样可以把内存拷贝降到最低。但对于大多数协议解析场景Buffer.concat 的简单性更值得优先考虑只有确认有性能瓶颈再上这种优化。4.4 手动 read(size) 的天然背压read(size) 除了精准还有一个隐藏优势它能帮你做节流。流动模式下数据到了你就得处理处理慢数据就会被积压。暂停模式下你不 read数据就堆在内部缓冲区堆满 highWaterMark 后底层就不再拉数据。这个特性在很多场景下非常好用。比如对接一个高速数据源平均每秒产生 10MB 数据而你的业务逻辑每秒只能处理 2MB。如果放任 data 事件流动内存会在几秒内被撑爆。而用 readable read(size)你可以实现一个最简单的处理完再读下一个的节流循环let paused true; stream.on(readable, () { if (!paused) return; paused true; while (true) { const chunk stream.read(1024); if (chunk null) break; // 这里如果处理特别耗时可以改成异步 循环控制 processChunk(chunk); } paused false; });当然真正的背压还得配合异步处理单纯同步读并不能解决异步处理慢的问题。但至少read(size) 给你提供了不读就不消费这种控制权这是 data 事件给不了的。5. 从限速到传感器读重read(size) 的应用扩展5.1 用 read(size) 做最朴素的限速说到限速很多人第一反应是引入复杂的流控库。其实利用 read(size) 加暂停模式就能做一个很简洁的按固定字节数限速的逻辑。思路是每次只从流里读出一小块数据处理完等一会儿再读下一块。const { Readable } require(stream); async function slowRead(stream, chunkSize 1024, delayMs 200) { let reading false; stream.on(readable, () { if (reading) return; reading true; (async () { while (true) { const chunk stream.read(chunkSize); if (chunk null) break; await new Promise(resolve setTimeout(resolve, delayMs)); processChunk(chunk); } reading false; })(); }); }这个函数模拟的是每 200ms 只消费 1024 字节的节奏非常适合那些下游设备能力弱、不能接受全速推送的场景。第一次 readable 事件触发后异步 IIFE 开始执行不断 read(1024)每次读完后 sleep 一段时间。读不到数据时退出循环把 reading 标志复位等待下次 readable 事件再启动新的读一批循环。这样流内部缓冲区既不会堆积太多数据下游也几乎不会感觉到压力。5.2 串口读取秤重量read(size) 的典型应用现在回到你大概率关心的场景用 Node.js 读取电子秤的重量。工业秤和传感器通常通过串口输出数据格式按帧算。以市面上常见的上海耀华、托利多等品牌为例虽然协议细节不同但共性非常明显每一帧固定长度以帧头开头中间是重量数据最后是校验。这正是 read(size) 发挥优势的地方。假设你的秤输出协议是每帧 18 字节第 0 字节是帧头 0xAA第 1 字节是帧头 0x55第 2~5 字节是有符号 32 位整数重量值单位克。使用 serialport 库接收串口数据再用 Readable 包装一层const { SerialPort } require(serialport); const { Readable } require(stream); const port new SerialPort({ path: COM3, // Windows 下串口Linux 下可能是 /dev/ttyUSB0 baudRate: 9600, }); const source new Readable({ read() {}, highWaterMark: 1024, }); port.on(data, (data) { // 串口数据到达推入自定义 Readable 的内部缓冲 source.push(data); }); source.on(readable, () { let frame; while ((frame source.read(18)) ! null) { const header1 frame[0]; const header2 frame[1]; if (header1 ! 0xAA || header2 ! 0x55) { console.error(帧头错误数据流可能失步); continue; } const weight frame.readInt32LE(2); console.log(当前重量${weight} 克); } });这个方案最大的优势在于你不需要关心串口每次 data 事件吐出来的数据是 3 字节还是 32 字节source.read(18) 会自动帮你从内部缓冲里跨越多个 chunk 拼出完整的帧。我第一次做这个需求时还在傻傻地自己写拼包逻辑后来发现 Readable 内部的 BufferList 已经帮你把这个脏活干完了。5.3 大文件分段处理read(size) 控制内存占用还有一个容易忽略的场景是读大文件。很多人喜欢用 fs.createReadStream 配合 data 事件处理但在处理超大文件时一次 data 事件可能给你 64KB 数据你要是把这些数据全部缓存到数组里再统一处理内存会非常难看。用 read(size) 可以精确控制每一批处理多少字节比如处理日志文件时一次读 4096 字节解析完再读下一批const fs require(fs); const stream fs.createReadStream(/path/to/huge.log, { highWaterMark: 8192, }); stream.on(readable, () { let chunk; while ((chunk stream.read(4096)) ! null) { // 按 4KB 为一批处理日志 processBatch(chunk.toString(utf8)); } });需要注意一点read(size) 不是保证一定返回 size 字节在流末尾时它可能返回少于 size 的字节。比如文件总共 12345 字节你每次都 read(4096)最后一次只会拿到 345 字节而不是 null。这个行为是符合预期的但你处理时要对最后一次可能不是整块有心理准备。5.4 结合 Transform 做协议重组read(size) 适合在底层做字节级控制但在一个复杂系统里你不可能把 read(18) 散落到业务代码的各个角落。更好的实践是把 read(size) 的拆帧逻辑封装在一个自定义 Transform 流里对外只输出完整的帧对象。const { Transform } require(stream); class FrameTransform extends Transform { constructor(frameLength) { super({ readableObjectMode: false }); this.frameLength frameLength; this.buffer Buffer.alloc(0); } _transform(chunk, encoding, callback) { this.buffer Buffer.concat([this.buffer, chunk]); while (this.buffer.length this.frameLength) { const frame this.buffer.subarray(0, this.frameLength); this.buffer this.buffer.subarray(this.frameLength); this.push(frame); } callback(); } _flush(callback) { if (this.buffer.length 0) { this.push(this.buffer); } callback(); } } const parser new FrameTransform(18); source.pipe(parser).on(data, (frame) { // 这里拿到的 frame 一定是完整的 18 字节 handleFrame(frame); });熟悉流处理的人会发现Transform 方案虽然本质上没有直接用 read(size)但它的思路完全一致把数据流切成定长的单位再交付给业务。用 read(size) 做底层协议解析用 Transform 做上层业务封装这两者是互补关系不是互斥关系。6. 到底该用 read(size) 还是交给 data/pipe我的选择原则6.1 四种消费方式的适用场景看到这里你可能会想那以后是不是所有场景都用 readable read(size) 就对了并不是。数据消费方式没有银弹我根据自己的经验整理了一个选择参考方式核心特点推荐场景data 事件被动接收代码最简单对 chunk 边界无严格要求、整体处理即可的数readable read()手动读全部缓冲控制权更强需要暂停/恢复、需要自己决定何时消费数据的场景readable read(size)精确按字节数读取二进制协议解析、定长帧拆包、限速、传感器设备对接pipe / pipeline自动背压、数据流转文件复制、数据透传、多个 Transform 串联处理划分的核心依据只有一条你的消费单位是什么。如果你的消费单位是文件整体或者一大块文本用 pipe 或 data 最省心如果你的消费单位是18 字节一帧或者每 1024 字节做一个快照那必须用 read(size)没有商量余地。6.2 我自己的几个经验总结最后分享一点我这个老油条的心得。第一次接触 read(size) 时我觉得这 API 很别扭不像 data 事件那样直观。但真正用习惯之后会发现它把数据边界这个复杂问题从业务代码里剥离了。有几个细节值得刻在脑子里。第一凡是二进制协议解析一定不要 setEncoding那会让 size 单位变成字符数。第二read(size) 返回 null 不代表结束只是当前不够要配合 readable 事件循环读。第三帧数据要跨异步边界处理时先 copy 一份 Buffer别图省事直接引用。第四end 事件里记得 read() 一次把残留数据捞出来否则最后那半个帧会悄无声息地消失。我之前做串口秤项目时最大的收益不是代码变短了而是排查问题的难度降了一个数量级。以前用 data 事件拆包粘包和半包问题总是偶发、难复现调试起来头大改用 read(size) 之后帧边界变得特别稳定几乎没再出现过数据错位的问题。如果你现在正在跟协议解析、串口数据、或者任何数据是一帧一帧到达的场景较劲认真尝试一下 read(size)大概率会回来感谢这个 API。
返回列表