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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 05:42:53