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

如何在NestJS端点实现流式传输并解决连接重置问题

解决MongoDB数据流流式传输时curl连接重置的问题

问题描述

我正在开发一个微服务,需在收到请求时向核心服务流式传输数据。最初的代码能处理大JSON负载,但会返回100个超大JSON对象,不够优化:

@Get('/zzz/stream')
async streamTiles(
  @Param('filterId', new ParseIntPipe()) filterId: number,
  @Query() query: Record<string, string> = {},
  @Res() res: Response
) {
  //... 省略stream获取逻辑,stream为MongoDB数据流

  await new Promise((resolve, reject) => {
    stream
      .pipe(JSONStream.stringify())
      .pipe(res)
      .on('finish', resolve)
      .on('error', reject);
  });

为实现客户端侧的轻量流式响应,我修改了代码,但收到错误:curl: (56) Recv failure: Connection reset by peer,修改后的代码如下:

@Get('/zzz/stream')
async streamTiles(
  @Param('filterId', new ParseIntPipe()) filterId: number,
  @Query() query: Record<string, string> = {},
  @Res({ passthrough: true }) res: Response
) {
   //... 省略stream获取逻辑,stream为MongoDB数据流

   res.setHeader('Content-Type', 'application/json');
   res.setHeader('Transfer-Encoding', 'chunked');
   res.setHeader('X-Content-Type-Options', 'nosniff');

   stream.pipe(res);

已知stream是有效的MongoDB游标数据流,需解决curl返回空响应及连接重置的问题。

解决方案

1. 正确序列化MongoDB流数据

MongoDB游标流输出原始BSON对象,直接pipe到响应会生成非合法JSON,导致客户端解析失败、连接重置。需逐个序列化文档,选择合适的流式JSON格式:

方案A:使用NDJSON(换行分隔JSON)

每个JSON对象单独一行,客户端可逐行解析,无需拼接完整数组:

@Get('/zzz/stream')
async streamTiles(
  @Param('filterId', new ParseIntPipe()) filterId: number,
  @Query() query: Record<string, string> = {},
  @Res({ passthrough: true }) res: Response
) {
  // 省略stream获取逻辑
  res.setHeader('Content-Type', 'application/x-ndjson');
  res.setHeader('X-Content-Type-Options', 'nosniff');

  stream.on('data', (doc) => {
    res.write(JSON.stringify(doc) + '\n');
  });

  stream.on('end', () => res.end());
  stream.on('error', (err) => res.status(500).end(JSON.stringify({ error: err.message })));
}

测试命令:curl http://your-service/zzz/stream,将看到逐行输出的JSON对象。

方案B:流式输出JSON数组

若需保持JSON数组格式,需手动处理数组的开头、元素分隔符和结尾:

@Get('/zzz/stream')
async streamTiles(
  @Param('filterId', new ParseIntPipe()) filterId: number,
  @Query() query: Record<string, string> = {},
  @Res({ passthrough: true }) res: Response
) {
  // 省略stream获取逻辑
  res.setHeader('Content-Type', 'application/json');
  res.setHeader('X-Content-Type-Options', 'nosniff');

  let isFirstDoc = true;
  res.write('[');

  stream.on('data', (doc) => {
    if (!isFirstDoc) res.write(',');
    isFirstDoc = false;
    res.write(JSON.stringify(doc));
  });

  stream.on('end', () => {
    res.write(']');
    res.end();
  });

  stream.on('error', (err) => res.status(500).end(JSON.stringify({ error: err.message })));
}

2. 移除手动设置的分块传输编码

Express/NestJS会自动处理Transfer-Encoding: chunked,手动设置可能引发冲突,直接删除res.setHeader('Transfer-Encoding', 'chunked');即可。

3. 确保响应生命周期正确结束

使用@Res({ passthrough: true })时,必须监听stream的end事件并调用res.end(),否则连接会挂起直至超时重置——这是本次连接重置的核心原因之一。

4. 完善错误处理逻辑

必须监听stream的error事件,否则MongoDB流抛出异常时会直接终止进程,导致连接被强制关闭。

内容的提问来源于stack exchange,提问作者Oliver Kucharzewski

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 09:15:23