NodeJS Stream提前终止,仅降低highWaterMark值可正常运行
Node.js Stream 提前终止问题修复
问题根源
- 异步生成器
processStream中未等待processChunk的Promise解析,直接yield了Promise对象,导致后续Transform流接收到的不是处理后的字符串,而是Promise,既造成数据异常,也破坏了流的背压机制。 - Transform流的
transform方法未正确处理res.write的异步特性,当响应缓冲区已满时直接调用callback(),会导致数据丢失或流提前终止。 - 额外监听
transformStream的end事件并调用res.end(),与pipeline的自动结束逻辑冲突,可能导致响应提前关闭。
修复方案
1. 等待异步Chunk处理完成
修改processStream,确保processChunk的Promise解析后再yield结果:
async function* processStream(source, { signal }) { source.setEncoding('utf-8'); for await (const chunk of source) { yield await processChunk(chunk); } }
2. 正确处理响应写入的背压
在Transform流的transform方法中,根据res.write的返回值判断是否需要等待drain事件:
const transformStream = new Transform({ async transform(chunk, encoding, callback) { const s = chunk.toString(); const canWrite = res.write(`data: ${s}\n\n`); if (!canWrite) { res.once('drain', callback); } else { callback(); } }, });
3. 统一响应结束时机
移除transformStream的end事件监听,在pipeline执行完成后再结束响应:
await pipeline(readStream, processStream, transformStream); res.end();
完整修复后代码
const express = require('express'); const app = express(); const port = 3010; const path = require('path'); const { Transform } = require('node:stream'); const { pipeline } = require('node:stream/promises'); const fs = require('node:fs'); app.use(express.static('static')); app.get('/', (req, res) => { res.sendFile(path.resolve('pages/index.html')); }); const processChunk = async function (chunk) { // simulate delay await new Promise((resolve, reject) => { setTimeout(resolve, 50); }); return chunk.toUpperCase(); }; app.get('/hello', async (req, res) => { res.type('text/event-stream'); res.setHeader('Cache-Control', 'no-cache'); const readStream = fs.createReadStream('./sample.txt', { highWaterMark: 5 * 1024, // 现在可以正常工作 }); async function* processStream(source, { signal }) { source.setEncoding('utf-8'); for await (const chunk of source) { yield await processChunk(chunk); } } const transformStream = new Transform({ async transform(chunk, encoding, callback) { const s = chunk.toString(); const canWrite = res.write(`data: ${s}\n\n`); if (!canWrite) { res.once('drain', callback); } else { callback(); } }, }); await pipeline(readStream, processStream, transformStream); res.end(); }); app.listen(port, () => { console.log(`Example app listening at http://localhost:${port}`); });
内容的提问来源于stack exchange,提问作者painotpi
相关产品推荐
相关产品推荐

