You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 00:53:23