在构建现代 Web 应用,特别是集成大模型(LLM)的智能应用时,我们常常会采用 BFF(Backend for Frontend)架构。Node.js 凭借其非阻塞 I/O 和事件驱动特性,成为实现 BFF 层的绝佳选择。当 BFF 需要将大模型(如 OpenAI GPT、Claude 等)的流式响应(Server-Sent Events, SSE)转发给前端时,一个关键但容易被忽视的问题浮出水面:如果用户在生成过程中关闭了浏览器标签页,或者网络突然中断,BFF 层与大模型后端的连接以及相关的计算资源该如何优雅地释放?

如果不妥善处理,这些“僵尸”连接和未完成的请求会持续消耗服务器资源(如内存、CPU、网络连接),最终可能导致服务内存泄漏、连接数耗尽,甚至拖垮整个应用。本文将深入剖析这一场景,从原理到实践,提供一套完整的 Node.js BFF 层处理 SSE 流式响应及客户端意外断开时资源释放的解决方案。无论你是正在搭建 AI 应用的中高级开发者,还是希望优化现有服务稳定性的工程师,都能从中获得可直接复用的代码和清晰的排查思路。

1. 背景与核心概念:为什么资源释放如此重要?

在深入代码之前,我们有必要厘清几个核心概念,并理解问题产生的根源。

1.1 BFF 层与流式响应

BFF(Backend for Frontend) :并非一个具体的技术,而是一种架构模式。它在前端与多个后端微服务(如用户服务、订单服务、大模型服务)之间,构建一个专为特定前端(如 Web 应用、移动端)定制的后端层。BFF 的核心职责是聚合、裁剪和转换后端服务的数据,为前端提供“恰好所需”的 API,从而简化前端逻辑,提升用户体验。

SSE(Server-Sent Events) :一种允许服务器通过 HTTP 连接主动向客户端推送数据的技术。与 WebSocket 的双向通信不同,SSE 是单向的(服务器到客户端),基于简单的文本协议。它非常适合实时通知、新闻推送、以及我们重点讨论的 大模型流式文本生成 。客户端使用 EventSource API 进行连接和监听。

流式响应(Streaming Response) :大模型(LLM)在处理复杂问题时,生成完整答案可能需要数秒甚至数十秒。为了提供更即时的反馈,服务端可以采用“流式”输出,即生成一个词或一段话就立刻发送给客户端,而不是等待全部生成完毕再一次性返回。这极大地降低了用户感知的延迟。

1.2 问题场景:客户端意外断开

在一个典型的 AI 问答应用中,流程如下:

  1. 用户在网页输入问题,点击“发送”。
  2. 前端通过 EventSource fetch API 向 Node.js BFF 发起一个 SSE 请求。
  3. Node.js BFF 收到请求后,作为代理,向真正的大模型服务(如 OpenAI API)发起另一个流式 HTTP 请求。
  4. 大模型服务开始流式返回数据块(chunks)。
  5. Node.js BFF 收到每个数据块后,立即将其按照 SSE 格式( data: ...\n\n )转发给前端。
  6. 前端通过 EventSource onmessage 事件实时渲染这些数据块。

风险点出现在第 4、5 步之间 。如果此时用户关闭了浏览器标签页,或者网络连接不稳定导致 TCP 连接断开,前端到 BFF 的连接会中断。然而,BFF 到后端大模型服务的 HTTP 流式请求可能仍在进行中。Node.js 默认不会自动终止这个下游请求,导致:

  • 内存泄漏 :持续接收的流数据会堆积在缓冲区。
  • 连接泄漏 :一个 HTTP(S) 连接被长期占用,无法被复用或关闭。
  • 不必要的计算成本 :大模型服务仍在为已离开的用户消耗宝贵的 GPU/CPU 资源。
  • 潜在的服务雪崩 :在高并发场景下,大量此类“僵尸请求”会快速耗尽服务器的文件描述符、内存和线程池资源,导致新请求无法处理。

因此, 在 BFF 层及时检测客户端断开并主动释放上下游资源,是保障服务健壮性的关键

2. 环境准备与版本说明

本文将使用 Node.js 和 Express 框架来构建 BFF 层示例。为了模拟大模型服务,我们会创建一个简单的模拟流式服务,并使用 axios node-fetch 来发起下游请求。

推荐环境:

  • 操作系统 :macOS / Linux / Windows (WSL2 推荐)
  • Node.js :>= 18.x (本文示例使用 Node.js 20+,因其对 fetch API 有稳定的原生支持)
  • 包管理器 :npm 或 yarn
  • IDE :VS Code 或其他你熟悉的编辑器

项目初始化: 首先,创建一个新的项目目录并初始化。

mkdir nodejs-bff-sse-cleanup
cd nodejs-bff-sse-cleanup
npm init -y

安装核心依赖: 我们将使用 Express 作为 Web 框架,并使用原生的 fetch (Node.js 18+ 内置)进行下游请求。为了更好的请求控制,我们也会介绍 axios 的用法。

npm install express
# 如果使用 Node.js < 18,或者希望使用功能更丰富的 HTTP 客户端,可以安装 axios
# npm install axios

项目结构预览:

nodejs-bff-sse-cleanup/
├── package.json
├── server.js              # 主 BFF 服务器文件
├── mock-llm-server.js     # 模拟的大模型流式服务
└── client.html            # 一个简单的前端测试页面

3. 核心原理与 Node.js 事件侦听

要解决问题,必须先理解 Node.js 的 http.ServerResponse 对象和流(Stream)的生命周期事件。

3.1 检测客户端连接关闭

在 Node.js 的 HTTP 服务器中,当客户端断开连接时,底层的 socket 会触发 'close' 'end' 事件。对于 SSE 这种长连接,我们需要在响应对象 ( res ) 上监听这些事件。

关键事件:

  • res.on('close', callback) :当底层连接被提前终止(例如客户端关闭标签页)时触发。这是 最可靠 的客户端断开检测信号。
  • res.on('finish', callback) :当响应已被完全发送(即所有数据已刷新到网络)时触发。在 SSE 场景下,连接是持久的,通常不会触发 finish ,除非你主动结束响应。
  • res.socket :通过 res.socket 可以访问到原始的 TCP socket,监听其 'close' 事件也能达到类似效果。

3.2 中止下游 Fetch 请求

从 Node.js 18 开始,原生的 fetch API 提供了 AbortController 来中止请求。这是释放下游资源的核心机制。

const controller = new AbortController();
const signal = controller.signal;

fetch('https://api.openai.com/v1/chat/completions', {
  method: 'POST',
  headers: { /* ... */ },
  body: JSON.stringify({ /* ... */ }),
  signal: signal // 传入 abort signal
})
  .then(response => { /* ... */ })
  .catch(err => {
    if (err.name === 'AbortError') {
      console.log('Fetch request was aborted');
    }
  });

// 当需要中止请求时
controller.abort();

调用 controller.abort() 会立即终止与该 fetch 请求相关的所有网络活动和流处理。

3.3 处理可读流(ReadableStream)

大模型服务返回的响应体是一个可读流。我们需要持续地从该流中读取数据,并转发给客户端。同时,我们必须确保在客户端断开时,停止读取并销毁这个流。

const downstreamResponse = await fetch(/* ... */, { signal });
const reader = downstreamResponse.body.getReader(); // 获取流阅读器

try {
  while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    // 将 value (Uint8Array) 转换为字符串并转发给前端 SSE
    res.write(`data: ${new TextDecoder().decode(value)}\n\n`);
  }
} catch (err) {
  // 处理错误,包括因 abort 导致的错误
} finally {
  reader.releaseLock(); // 重要:释放阅读器锁
}

4. 完整实战案例:构建健壮的 BFF SSE 代理

让我们一步步构建一个完整的、具备资源释放能力的 BFF 服务。

4.1 创建模拟的大模型流式服务

首先,我们创建一个独立的模拟服务 ( mock-llm-server.js ),它模拟 OpenAI 等服务的流式响应。这有助于我们在本地进行测试,而无需消耗真实的 API 额度。

// file: mock-llm-server.js
const http = require('http');

const server = http.createServer((req, res) => {
  // 只处理特定路径的 POST 请求
  if (req.url === '/v1/chat/completions' && req.method === 'POST') {
    console.log(`[Mock LLM] Received request from ${req.socket.remoteAddress}`);

    // 设置 SSE 响应头
    res.writeHead(200, {
      'Content-Type': 'text/event-stream',
      'Cache-Control': 'no-cache',
      'Connection': 'keep-alive',
      'Access-Control-Allow-Origin': '*', // 为了方便测试,允许跨域
    });

    // 模拟一个长时间的流式响应
    const message = "这是一个来自模拟大模型的流式响应。它将被分成多个块发送。";
    const chunks = message.split(''); // 按字符分割,模拟 token 流
    let index = 0;

    const intervalId = setInterval(() => {
      if (index >= chunks.length) {
        // 发送结束标志
        res.write('data: [DONE]\n\n');
        clearInterval(intervalId);
        res.end(); // 结束响应
        console.log(`[Mock LLM] Stream finished for ${req.socket.remoteAddress}`);
        return;
      }

      const chunk = chunks[index];
      // 模拟 SSE 数据格式,通常大模型 API 返回的是 JSON 字符串
      const data = JSON.stringify({ choices: [{ delta: { content: chunk } }] });
      res.write(`data: ${data}\n\n`);
      console.log(`[Mock LLM] Sent chunk: ${chunk}`);
      index++;
    }, 100); // 每 100ms 发送一个字符

    // 关键:监听客户端断开连接
    req.on('close', () => {
      console.log(`[Mock LLM] Client disconnected from ${req.socket.remoteAddress}. Stopping stream.`);
      clearInterval(intervalId); // 停止发送
      // 在实际的大模型服务中,这里应该通知模型停止生成
    });

    // 监听请求体结束(可选,用于获取请求数据)
    let body = '';
    req.on('data', chunk => { body += chunk; });
    req.on('end', () => {
      console.log(`[Mock LLM] Request body: ${body}`);
    });

  } else {
    res.writeHead(404);
    res.end('Not Found');
  }
});

const PORT = 3001;
server.listen(PORT, () => {
  console.log(`Mock LLM Server running at http://localhost:${PORT}`);
});

运行 node mock-llm-server.js 启动模拟服务。

4.2 构建 Node.js BFF 服务器(基础版,无清理)

我们先写一个基础的 BFF,它只是简单地代理请求,但 没有 处理客户端断开。

// file: server-basic.js
const express = require('express');
const app = express();
app.use(express.json());

app.post('/api/chat/stream', async (req, res) => {
  console.log(`[BFF] Received request from ${req.ip}`);

  // 1. 设置 SSE 响应头
  res.writeHead(200, {
    'Content-Type': 'text/event-stream',
    'Cache-Control': 'no-cache',
    'Connection': 'keep-alive',
    'Access-Control-Allow-Origin': '*', // 生产环境应严格限制
  });

  // 2. 转发请求到下游大模型服务
  const downstreamResponse = await fetch('http://localhost:3001/v1/chat/completions', {
    method: 'POST',
    headers: {
      'Content-Type': 'application/json',
      // 这里可以添加认证头,如 'Authorization': `Bearer ${API_KEY}`
    },
    body: JSON.stringify(req.body), // 将前端请求体原样转发
  });

  // 3. 获取下游的流式响应体
  const reader = downstreamResponse.body.getReader();
  const decoder = new TextDecoder();

  try {
    while (true) {
      const { done, value } = await reader.read();
      if (done) {
        res.write('data: [DONE]\n\n');
        res.end();
        console.log(`[BFF] Downstream stream finished for ${req.ip}`);
        break;
      }
      // 4. 将下游的数据块转发给前端
      const chunk = decoder.decode(value);
      res.write(`data: ${chunk}\n\n`);
      console.log(`[BFF] Forwarded chunk: ${chunk.substring(0, 50)}...`);
    }
  } catch (error) {
    console.error(`[BFF] Error reading stream for ${req.ip}:`, error.message);
    if (!res.headersSent) {
      res.writeHead(500);
    }
    res.end();
  }
});

const PORT = 3000;
app.listen(PORT, () => {
  console.log(`BFF Server (Basic) running at http://localhost:${PORT}`);
  console.log(`Test endpoint: POST http://localhost:${PORT}/api/chat/stream`);
});

启动这个 BFF ( node server-basic.js ) 并用工具测试,你会发现如果客户端中途断开,BFF 控制台会继续打印日志,直到模拟 LLM 服务发送完所有数据。这证明了资源泄漏正在发生。

4.3 构建 Node.js BFF 服务器(增强版,带资源释放)

现在,我们加入资源释放的核心逻辑。

// file: server-advanced.js
const express = require('express');
const app = express();
app.use(express.json());

app.post('/api/chat/stream', async (req, res) => {
  console.log(`[BFF] Received request from ${req.ip}`);

  // 1. 设置 SSE 响应头
  res.writeHead(200, {
    'Content-Type': 'text/event-stream',
    'Cache-Control': 'no-cache',
    'Connection': 'keep-alive',
    'Access-Control-Allow-Origin': '*',
  });

  // 2. 创建 AbortController 用于控制下游请求
  const downstreamAbortController = new AbortController();
  const downstreamSignal = downstreamAbortController.signal;

  // 3. 监听客户端连接关闭事件 - 这是资源释放的触发器
  let isClientConnected = true;
  const cleanup = () => {
    if (!isClientConnected) return; // 防止重复清理
    isClientConnected = false;
    console.log(`[BFF] Client ${req.ip} disconnected. Cleaning up.`);

    // 中止下游的 fetch 请求
    downstreamAbortController.abort();
    // 注意:我们不需要手动调用 res.end(),因为连接已关闭,尝试写入会报错。
  };

  // 主要监听 'close' 事件
  res.on('close', cleanup);
  // 也可以监听 socket 的 close 事件作为额外保障
  req.socket.on('close', cleanup);

  // 4. 设置请求超时(可选但推荐)
  const requestTimeout = setTimeout(() => {
    console.log(`[BFF] Request timeout for ${req.ip}`);
    cleanup();
    if (!res.headersSent) {
      res.writeHead(408);
    }
    res.end();
  }, 60000); // 60秒超时

  try {
    // 5. 发起下游请求,传入 abort signal
    const downstreamResponse = await fetch('http://localhost:3001/v1/chat/completions', {
      method: 'POST',
      headers: {
        'Content-Type': 'application/json',
      },
      body: JSON.stringify(req.body),
      signal: downstreamSignal, // 关键:绑定 abort 信号
    });

    // 6. 检查下游响应状态
    if (!downstreamResponse.ok) {
      throw new Error(`Downstream error: ${downstreamResponse.status}`);
    }

    // 7. 获取流阅读器
    const reader = downstreamResponse.body.getReader();
    const decoder = new TextDecoder();

    // 8. 循环读取并转发数据
    try {
      while (isClientConnected) { // 循环条件检查客户端是否还在线
        const { done, value } = await reader.read();
        if (done) {
          console.log(`[BFF] Downstream stream finished normally for ${req.ip}`);
          res.write('data: [DONE]\n\n');
          break;
        }

        const chunk = decoder.decode(value, { stream: true });
        // 重要:在写入前再次检查连接状态
        if (!isClientConnected) {
          console.log(`[BFF] Client disconnected during write. Aborting.`);
          break;
        }

        // 转发数据
        res.write(`data: ${chunk}\n\n`);
        console.log(`[BFF] Forwarded chunk for ${req.ip}`);
      }
    } catch (streamError) {
      // 这个 catch 主要捕获 reader.read() 的异常,包括因 abort 产生的 AbortError
      if (streamError.name === 'AbortError') {
        console.log(`[BFF] Downstream stream reading was aborted for ${req.ip}`);
      } else {
        console.error(`[BFF] Error reading downstream stream for ${req.ip}:`, streamError);
      }
    } finally {
      // 9. 最终清理:释放阅读器锁,清除超时定时器
      reader.releaseLock();
      clearTimeout(requestTimeout);
      // 如果客户端还连着,正常结束响应;如果已经断了,res.end() 可能会报错,忽略即可。
      if (isClientConnected) {
        res.end();
      }
      console.log(`[BFF] Request processing finished for ${req.ip}`);
    }

  } catch (fetchError) {
    // 捕获 fetch 本身的错误(如网络错误、因 abort 导致的错误)
    clearTimeout(requestTimeout);
    if (fetchError.name === 'AbortError') {
      console.log(`[BFF] Downstream fetch was aborted for ${req.ip}`);
    } else {
      console.error(`[BFF] Fetch error for ${req.ip}:`, fetchError.message);
      if (isClientConnected && !res.headersSent) {
        res.writeHead(502);
        res.end(JSON.stringify({ error: 'Failed to connect to upstream service' }));
      }
    }
  }
});

const PORT = 3000;
app.listen(PORT, () => {
  console.log(`BFF Server (Advanced) running at http://localhost:${PORT}`);
  console.log(`Test endpoint: POST http://localhost:${PORT}/api/chat/stream`);
});

4.4 创建前端测试页面

创建一个简单的 HTML 页面来测试我们的 BFF。

<!-- file: client.html -->
<!DOCTYPE html>
<html lang="en">
<head>
    <meta charset="UTF-8">
    <meta name="viewport" content="width=device-width, initial-scale=1.0">
    <title>SSE Client Test</title>
</head>
<body>
    <h1>SSE Stream Test with BFF</h1>
    <button id="startBtn">Start Stream</button>
    <button id="stopBtn" disabled>Stop Stream (Close Connection)</button>
    <button id="abortBtn" disabled>Abort Fetch (模拟网络错误)</button>
    <br><br>
    <div id="output" style="white-space: pre-wrap; border:1px solid #ccc; padding:10px; min-height:200px;"></div>

    <script>
        let eventSource = null;
        let fetchController = null;
        const output = document.getElementById('output');

        document.getElementById('startBtn').onclick = async () => {
            output.textContent = 'Starting stream...\n';
            document.getElementById('startBtn').disabled = true;
            document.getElementById('stopBtn').disabled = false;
            document.getElementById('abortBtn').disabled = false;

            // 方法1: 使用 EventSource (标准 SSE)
            // eventSource = new EventSource('http://localhost:3000/api/chat/stream');
            // eventSource.onmessage = (e) => {
            //     output.textContent += `[EventSource] ${e.data}\n`;
            // };
            // eventSource.onerror = (e) => {
            //     output.textContent += `[EventSource Error] Connection closed.\n`;
            //     cleanup();
            // };

            // 方法2: 使用 Fetch API 读取流 (更灵活,可以发送 POST 和 Body)
            fetchController = new AbortController();
            try {
                const response = await fetch('http://localhost:3000/api/chat/stream', {
                    method: 'POST',
                    headers: { 'Content-Type': 'application/json' },
                    body: JSON.stringify({ messages: [{ role: 'user', content: 'Hello' }] }),
                    signal: fetchController.signal
                });

                const reader = response.body.getReader();
                const decoder = new TextDecoder();

                while (true) {
                    const { done, value } = await reader.read();
                    if (done) {
                        output.textContent += '[Fetch] Stream finished.\n';
                        break;
                    }
                    // SSE 数据格式是 "data: {...}\n\n",需要解析
                    const text = decoder.decode(value);
                    const lines = text.split('\n');
                    for (const line of lines) {
                        if (line.startsWith('data: ')) {
                            const data = line.substring(6); // 去掉 "data: "
                            if (data.trim() === '[DONE]') {
                                output.textContent += '[Fetch] Stream completed with [DONE].\n';
                            } else {
                                try {
                                    const parsed = JSON.parse(data);
                                    const content = parsed.choices?.[0]?.delta?.content || '';
                                    if (content) {
                                        output.textContent += content;
                                    }
                                } catch(e) {
                                    output.textContent += `[Raw Data] ${data}\n`;
                                }
                            }
                        }
                    }
                }
            } catch (err) {
                if (err.name === 'AbortError') {
                    output.textContent += '[Fetch] Request was aborted by user.\n';
                } else {
                    output.textContent += `[Fetch Error] ${err.message}\n`;
                }
            } finally {
                cleanup();
            }
        };

        document.getElementById('stopBtn').onclick = () => {
            output.textContent += '[UI] Manually stopping connection.\n';
            cleanup();
        };

        document.getElementById('abortBtn').onclick = () => {
            if (fetchController) {
                output.textContent += '[UI] Manually aborting fetch request.\n';
                fetchController.abort();
            }
        };

        function cleanup() {
            if (eventSource) {
                eventSource.close();
                eventSource = null;
            }
            if (fetchController) {
                // 如果请求还在进行,abort 会触发 catch 块
                fetchController = null;
            }
            document.getElementById('startBtn').disabled = false;
            document.getElementById('stopBtn').disabled = true;
            document.getElementById('abortBtn').disabled = true;
        }
    </script>
</body>
</html>

4.5 运行与验证

  1. 启动服务 :打开三个终端窗口。
    • 终端1: node mock-llm-server.js (监听 3001 端口)
    • 终端2: node server-advanced.js (监听 3000 端口)
  2. 测试连接 :用浏览器打开 client.html 文件(可以直接双击,或通过 python3 -m http.server 8080 等服务访问)。
  3. 正常流程 :点击 "Start Stream",观察 BFF 和 Mock LLM 服务器的控制台输出,以及网页上逐渐出现的文字。
  4. 测试资源释放
    • 场景 A (关闭连接) :在流式输出过程中,点击 "Stop Stream" 按钮或直接关闭浏览器标签页。观察 BFF 控制台是否立即打印 Client ... disconnected. Cleaning up. Downstream fetch was aborted ,同时 Mock LLM 服务器是否打印 Client disconnected ... Stopping stream. 。这表明资源已被正确释放。
    • 场景 B (超时) :你可以将 BFF 代码中的 requestTimeout 改短(如 5000 毫秒)来测试超时逻辑。
    • 场景 C (网络错误) :点击 "Abort Fetch" 按钮,模拟前端主动中止请求。

通过对比 server-basic.js server-advanced.js 在客户端断开后的行为,你可以清晰地看到资源释放机制带来的差异。

5. 常见问题与排查思路

在实际部署中,你可能会遇到以下问题。这里提供一个排查清单。

问题现象 可能原因 排查步骤与解决方案
BFF 服务器内存使用量持续增长 1. 客户端断开连接未正确检测。
2. 下游流未被正确销毁。
3. 事件监听器未移除导致内存泄漏。
1. 确保 res.on('close', ...) 被正确绑定且能触发。
2. 使用 AbortController 并确认下游请求被中止。
3. 在清理函数中,移除不必要的监听器(如 req.socket.off('close', cleanup) )。
4. 使用 Node.js 内存分析工具(如 --inspect 配合 Chrome DevTools)查找泄漏点。
下游大模型服务连接未终止 1. AbortController.signal 未正确传递给 fetch 或 HTTP 客户端库。
2. 下游服务不支持请求中止。
1. 检查 fetch axios 配置,确保 signal 参数已设置。
2. 对于不支持 signal 的库(如 request ),考虑更换或手动销毁 socket。
3. 确认下游服务(如 OpenAI API)是否支持并正确处理 Connection: close 或请求中断。
res.write() 抛出 ECONNRESET 错误 在客户端已断开连接后,仍尝试向响应流写入数据。 1. 在每次 res.write() 前,检查 isClientConnected 标志位。
2. 使用 try...catch 包裹 res.write() 调用,并忽略 ECONNRESET 等特定错误。
3. 使用 res.writableEnded res.writableFinished 属性判断流是否可写(Node.js 版本需支持)。
客户端断开检测延迟或不触发 1. 网络环境复杂(如负载均衡、代理)。
2. 客户端非正常关闭(如断电、进程崩溃)。
1. 增加心跳机制:BFF 定期向前端发送注释行( :\n\n ),如果多次发送失败,可视为连接已死。
2. 结合 requestTimeout 设置一个全局请求超时,作为最后的保障。
3. 考虑在负载均衡器(如 Nginx)层面设置 proxy_read_timeout 并确保其能正确传递断开信号。
使用 axios 时流式响应处理异常 axios 的响应默认不是流,会缓冲整个响应体。 1. 在请求配置中设置 responseType: 'stream'
2. 使用 axios.CancelToken AbortController (axios >= 0.22.0)来取消请求。
3. 正确处理 axios 返回的 Node.js 流对象。

6. 最佳实践与工程建议

将上述解决方案投入生产环境,还需要考虑更多工程化细节。

6.1 使用中间件封装清理逻辑

将资源释放的逻辑抽象成 Express 中间件,可以提高代码复用性和可维护性。

// file: middleware/sseCleanup.js
const createSSEStreamHandler = (upstreamUrlFetcher, options = {}) => {
  return async (req, res, next) => {
    // 设置 SSE 头
    res.writeHead(200, {
      'Content-Type': 'text/event-stream',
      'Cache-Control': 'no-cache',
      'Connection': 'keep-alive',
    });

    const cleanupController = new AbortController();
    let isClientConnected = true;

    const cleanup = () => {
      if (!isClientConnected) return;
      isClientConnected = false;
      cleanupController.abort();
      // 移除监听器,防止内存泄漏
      res.off('close', cleanup);
      req.socket?.off('close', cleanup);
      if (timeoutId) clearTimeout(timeoutId);
      console.log(`[SSE Middleware] Cleanup completed for ${req.ip}`);
    };

    res.on('close', cleanup);
    req.socket?.on('close', cleanup);

    const timeoutMs = options.timeout || 60000;
    const timeoutId = setTimeout(() => {
      console.log(`[SSE Middleware] Timeout for ${req.ip}`);
      cleanup();
      if (!res.headersSent) {
        res.writeHead(408);
      }
      res.end();
    }, timeoutMs);

    try {
      // 调用外部函数获取上游 URL 和请求配置
      const { url, fetchOptions } = await upstreamUrlFetcher(req);
      const downstreamResponse = await fetch(url, {
        ...fetchOptions,
        signal: cleanupController.signal,
      });

      if (!downstreamResponse.ok) {
        throw new Error(`Upstream error: ${downstreamResponse.status}`);
      }

      const reader = downstreamResponse.body.getReader();
      const decoder = new TextDecoder();

      try {
        while (isClientConnected) {
          const { done, value } = await reader.read();
          if (done) break;
          if (!isClientConnected) break;
          const chunk = decoder.decode(value, { stream: true });
          // 可以在这里加入数据转换逻辑
          res.write(`data: ${chunk}\n\n`);
        }
        if (isClientConnected) {
          res.write('data: [DONE]\n\n');
        }
      } finally {
        reader.releaseLock();
        clearTimeout(timeoutId);
        if (isClientConnected) res.end();
      }
    } catch (error) {
      clearTimeout(timeoutId);
      if (error.name === 'AbortError') {
        // 预期内的中止,无需处理
      } else {
        console.error(`[SSE Middleware] Error for ${req.ip}:`, error);
        if (isClientConnected && !res.headersSent) {
          res.writeHead(502);
          res.end(JSON.stringify({ error: 'Stream failed' }));
        }
      }
    }
  };
};

module.exports = { createSSEStreamHandler };

在路由中使用:

// file: server-with-middleware.js
const express = require('express');
const { createSSEStreamHandler } = require('./middleware/sseCleanup');
const app = express();
app.use(express.json());

const chatStreamHandler = createSSEStreamHandler(async (req) => {
  // 根据业务逻辑动态构造上游请求
  return {
    url: 'http://localhost:3001/v1/chat/completions',
    fetchOptions: {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify(req.body),
    },
  };
}, { timeout: 120000 }); // 2分钟超时

app.post('/api/chat/stream', chatStreamHandler);

6.2 监控与日志

在生产环境中,完善的监控和日志至关重要。

  • 连接数监控 :监控 BFF 服务器的活跃 HTTP 连接数( server.getConnections() )和文件描述符使用量。
  • 错误日志 :记录所有连接断开、请求中止、下游服务错误的事件,并附上请求 ID 和客户端 IP,便于溯源。
  • 性能指标 :记录每个流式请求的持续时间、传输数据量,设置告警阈值。

6.3 考虑使用专门的流处理库

对于更复杂的场景(如背压处理、多路复用),可以考虑使用专门的流处理库,如:

  • pump / pipeline (Node.js 内置) :用于安全地管道化流,并自动处理错误和清理。
  • axios + stream :如前所述,正确配置。
  • got :一个功能强大的 HTTP 请求库,对流和取消有很好的支持。

6.4 安全与限流

  • 认证与授权 :在 BFF 层验证用户身份和权限,再决定是否转发请求到大模型服务。
  • 速率限制 :防止单个用户或 IP 过度消耗资源。可以使用 express-rate-limit 等中间件。
  • 请求体大小限制 :使用 express.json({ limit: '1mb' }) 限制请求体,防止内存耗尽攻击。
  • CORS :在生产环境中,将 Access-Control-Allow-Origin 设置为明确的前端域名,而不是 *

6.5 部署与运维

  • 进程管理 :使用 PM2、Docker 或 Kubernetes 管理 Node.js 进程,确保崩溃后能自动重启。
  • 优雅关闭 :在收到 SIGTERM 等信号时,应等待正在处理的流式请求完成或超时后再退出,避免数据丢失。
  • 负载均衡 :如果部署了多个 BFF 实例,确保 SSE 长连接会话的粘性(如果需要),或者确保客户端断开信号能正确传递。

7. 总结

在 Node.js BFF 层处理大模型 SSE 流式响应时,客户端的意外断开是一个必须严肃对待的工程问题。核心解决思路可以概括为 “侦听断开 -> 中止下游 -> 清理资源” 三步走。

  1. 可靠侦听 :通过 res.on('close', ...) req.socket 事件,及时感知客户端连接状态变化。
  2. 主动中止 :利用 AbortController fetch signal 参数,立即终止对下游大模型服务的请求,避免不必要的计算和网络开销。
  3. 全面清理 :释放流阅读器 ( reader.releaseLock() )、清除超时定时器、移除事件监听器,并将连接状态标志位 ( isClientConnected ) 纳入所有关键循环和写入操作的判断条件中。

本文提供的从基础到进阶的示例代码,以及中间件封装模式,为你构建生产级应用提供了可直接参考的蓝图。记住,资源释放不仅是代码正确性问题,更是服务稳定性和成本控制的关键。在实现核心功能后,务必结合监控、日志、安全限流和优雅关闭等工程实践,打造一个既高效又健壮的流式 AI 应用后端。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐