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

如何用AbortController实现Node.js流的暂停/恢复而非终止?

实现流的暂停与恢复功能

方案一:客户端本地控制暂停/恢复(无需服务端改动)

这种方式在客户端拦截流数据,暂停时缓存数据,恢复时再传递到目标节点,服务端仍会持续发送数据,适合小流量场景。

代码实现

// 全局保存流相关引用
let activeReadable = null;
let transformController = null;
let isPaused = false;
let cachedChunks = [];

// 启动流消费
start.addEventListener('click', async () => {
  // 避免重复启动
  if (activeReadable) return;

  // 获取服务端流(如需终止功能可保留AbortController)
  activeReadable = await consumeAPI();

  // 创建中间转换流,用于控制数据传递
  const transformStream = new TransformStream({
    transform(chunk, controller) {
      if (isPaused) {
        // 暂停时缓存数据
        cachedChunks.push(chunk);
      } else {
        // 正常传递数据
        controller.enqueue(chunk);
      }
    },
    flush(controller) {
      // 流结束时清空缓存
      cachedChunks.forEach(chunk => controller.enqueue(chunk));
      cachedChunks = [];
    }
  });

  transformController = transformStream.controller;

  // 管道连接:源流 -> 转换流 -> 目标DOM更新
  await activeReadable
    .pipeThrough(transformStream)
    .pipeTo(appendToHTML(cards));

  // 流结束后重置引用
  activeReadable = null;
  transformController = null;
});

// 暂停按钮逻辑
pause.addEventListener('click', () => {
  isPaused = true;
  console.log('流已暂停');
});

// 恢复按钮逻辑
resume.addEventListener('click', () => {
  if (!isPaused || !transformController) return;

  isPaused = false;
  // 发送缓存的所有数据
  cachedChunks.forEach(chunk => transformController.enqueue(chunk));
  cachedChunks = [];
  console.log('流已恢复');
});

方案二:服务端配合实现真正暂停/恢复(节省带宽)

如果需要让服务端停止发送数据(减少带宽占用),需要客户端给服务端发送控制信号,服务端调用Node.js Readable流的pause()/resume()方法。

服务端代码(Node.js)

const express = require('express');
const { Readable } = require('stream');
const { toWeb } = require('stream/web');

const app = express();
// 用Map管理多客户端流(示例简化为单客户端)
const clientStreams = new Map();
const clientId = 'demo-client';

// 提供流接口
app.get('/stream', (req, res) => {
  // 创建Node.js可读流
  const nodeReadable = new Readable({
    read(size) {
      // 模拟定时生成数据
      setTimeout(() => {
        if (!this.isPaused()) {
          this.push(`[${new Date().toLocaleTimeString()}] 数据块\n`);
        }
      }, 1000);
    }
  });

  clientStreams.set(clientId, nodeReadable);
  // 转换为Web Stream并响应
  const webReadable = toWeb(nodeReadable);
  res.setHeader('Content-Type', 'text/plain');
  webReadable.pipeTo(res);

  // 响应结束后清理流引用
  res.on('finish', () => {
    clientStreams.delete(clientId);
  });
});

// 暂停流接口
app.post('/pause-stream', (req, res) => {
  const stream = clientStreams.get(clientId);
  if (stream) {
    stream.pause();
    res.status(200).send('流已暂停');
  } else {
    res.status(404).send('无活跃流');
  }
});

// 恢复流接口
app.post('/resume-stream', (req, res) => {
  const stream = clientStreams.get(clientId);
  if (stream) {
    stream.resume();
    res.status(200).send('流已恢复');
  } else {
    res.status(404).send('无活跃流');
  }
});

app.listen(3000, () => console.log('服务启动在3000端口'));

客户端代码

let isStreamActive = false;

// 启动流消费
start.addEventListener('click', async () => {
  if (isStreamActive) return;
  isStreamActive = true;

  const response = await fetch('/stream');
  const readable = response.body;
  await readable.pipeTo(appendToHTML(cards));
  
  isStreamActive = false;
});

// 暂停按钮逻辑
pause.addEventListener('click', async () => {
  await fetch('/pause-stream', { method: 'POST' });
  console.log('已通知服务端暂停流');
});

// 恢复按钮逻辑
resume.addEventListener('click', async () => {
  await fetch('/resume-stream', { method: 'POST' });
  console.log('已通知服务端恢复流');
});

补充说明

  • 如果需要保留终止流的功能,可以继续保留AbortController,在终止时调用abort()并清理所有流引用。
  • 客户端本地缓存方案要注意内存占用,大流量场景建议使用服务端配合的方案。

内容的提问来源于stack exchange,提问作者Samuel A. Souza

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 21:30:49