如何在NestJS端点实现流式传输并解决连接重置问题
问题描述
我正在开发一个微服务,需在收到请求时向核心服务流式传输数据。最初的代码能处理大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

