如何在Node.js(NestJS)中检测客户端中止HTTP请求
解决NestJS SSE中客户端断开时终止服务端操作的方案
核心思路是利用Observable的订阅取消机制,结合NestJS/Express的连接断开事件,在客户端中止请求时触发服务端的资源清理逻辑,终止数据查询等异步操作。
方案一:通过Express响应的close事件触发取消信号
直接监听客户端连接断开事件,发送取消信号给Observable,终止内部异步任务。
Controller层实现
import { Controller, Sse, Res } from '@nestjs/common'; import { Response } from 'express'; import { Observable, Subject, finalize } from 'rxjs'; import { MessageEvent } from '@nestjs/common'; import { DataService } from './data.service'; @Controller('sse') export class SseController { constructor(private readonly dataService: DataService) {} @Sse('stream') sseStream(@Res() res: Response): Observable<MessageEvent> { // 创建取消信号的Subject const cancel$ = new Subject<void>(); // 监听客户端断开事件 res.on('close', () => { cancel$.next(); cancel$.complete(); }); return this.dataService.getStreamData(cancel$).pipe( finalize(() => { console.log('客户端已断开,完成资源清理'); }) ); } }
Service层实现
在Observable内部监听取消信号,终止异步任务(以定时查询为例,实际可替换为数据库查询/HTTP请求):
import { Injectable } from '@nestjs/common'; import { Observable, Subject, takeUntil } from 'rxjs'; import { MessageEvent } from '@nestjs/common'; @Injectable() export class DataService { getStreamData(cancel$: Subject<void>): Observable<MessageEvent> { return new Observable((subscriber) => { // 模拟持续的数据查询任务 const queryInterval = setInterval(async () => { // 实际场景:执行数据库查询/API请求 const data = await this.fetchDataFromDB(); subscriber.next({ data }); }, 1000); // 收到取消信号时终止任务 cancel$.subscribe(() => { clearInterval(queryInterval); subscriber.complete(); }); // 订阅取消时的兜底清理 return () => clearInterval(queryInterval); }).pipe(takeUntil(cancel$)); } private async fetchDataFromDB() { // 模拟数据库查询逻辑 return { id: Date.now(), content: '实时更新数据' }; } }
方案二:利用Observable自动取消的特性简化实现
NestJS处理SSE时,客户端断开连接会自动取消Observable的订阅,因此可以直接在Observable的清理函数中终止任务:
Controller层
@Sse('stream') sseStream(): Observable<MessageEvent> { return this.dataService.getStreamData().pipe( finalize(() => { console.log('客户端断开,停止所有异步操作'); }) ); }
Service层
getStreamData(): Observable<MessageEvent> { return new Observable((subscriber) => { const queryInterval = setInterval(async () => { const data = await this.fetchDataFromDB(); subscriber.next({ data }); }, 1000); // 订阅取消时自动执行此清理函数 return () => { clearInterval(queryInterval); console.log('终止数据查询任务'); }; }); }
处理可取消的异步任务(如HTTP/数据库查询)
对于支持取消的异步操作(如使用AbortController的HTTP请求、TypeORM流式查询),可以在清理函数中触发终止:
getStreamData(): Observable<MessageEvent> { return new Observable((subscriber) => { const abortController = new AbortController(); const signal = abortController.signal; const runQueryLoop = async () => { try { while (!signal.aborted) { // 带取消信号的HTTP请求 const response = await fetch('https://api.example.com/real-time-data', { signal }); const data = await response.json(); subscriber.next({ data }); await new Promise(resolve => setTimeout(resolve, 1000)); } } catch (err) { // 忽略取消导致的错误 if (err.name !== 'AbortError') { subscriber.error(err); } } finally { subscriber.complete(); } }; runQueryLoop(); // 订阅取消时终止请求 return () => abortController.abort(); }); }
内容的提问来源于stack exchange,提问作者Thai Nguyen
相关产品推荐
相关产品推荐

