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

Observable未取消订阅引发内存泄漏及Jest测试挂起问题

问题分析与解决方案

你的问题核心在于取消订阅信号没有正确传递到上游的obsWebsocket Observable,同时代码里存在几处关键的小错误,导致WebSocket连接无法被清理,retry逻辑也无法停止。让我们一步步拆解:

1. 修复obsWebsocket中的变量名错误

看你的obsWebsocket代码,ws.onclose回调里用了未定义的observer变量,正确的应该是参数里的obs:

ws.onclose = close => obs.error(new Error('close')); // 把observer改成obs

这个笔误会导致WebSocket关闭时抛出未捕获的ReferenceError,而不是把错误发送到Observable管道里。这种未捕获的错误会打乱RxJS的订阅清理逻辑,导致retry无法正确响应取消信号,WebSocket的清理函数也不会被触发。

另外,别忘了处理WebSocket的error事件——WebSocket的错误不会自动触发close事件,所以需要单独添加:

ws.onerror = (err) => obs.error(err);

2. 确保取消订阅信号传递到所有上游流

retry(3)操作符本身是会响应取消订阅的,但如果上游流的错误处理异常(比如上面的变量名错误),会导致这个逻辑失效。另外,你可以做一个优化:用takeUntil结合一个终止信号Subject来更可靠地取消所有上游订阅,避免依赖单一的subscription.unsubscribe():

// 在模块或测试文件中定义一个终止信号Subject
const cleanupSignal$ = new Subject<void>();

// 修改订阅逻辑
const subscription = getDataByTopic(topics).pipe(
  takeUntil(cleanupSignal$)
).subscribe(data => {
  // do stuff with data
});

cleanup() {
  // 发送终止信号,所有上游流都会收到取消信号
  cleanupSignal$.next();
  cleanupSignal$.complete();
  // 如果需要,也可以保留subscription.unsubscribe(),但takeUntil已经足够
}

这种方式比直接调用unsubscribe()更可靠,因为它能确保整个管道里的所有操作符(包括mergeMap、retry)都收到取消信号,进而触发上游obsWebsocket的清理函数。

3. 简化冗余的mergeMap

handleMsg里的mergeMap其实是多余的,因为of(parsed)是同步的单值流,换成map更简洁,也避免不必要的流创建:

function handleMsg(data: Observable<Buffer>): Observable<ParsedMessage> {
   return data.pipe(
      retry(3),
      map(msg => parseMsg(msg)) // 替换mergeMap为map
   )
}

4. Jest测试中的额外注意事项

在Jest端到端测试中,要确保cleanup()在每个测试结束时被调用,比如用afterEach钩子:

afterEach(() => {
  cleanup();
});

如果测试是异步的,还要确保Jest等待所有异步操作完成,比如使用async/await或者done回调,避免测试提前结束而没触发清理。

为什么取消订阅没触发清理?

正常情况下,取消主订阅会沿着RxJS管道向上传递,每个操作符都会取消对上游流的订阅,最终触发obsWebsocket的清理函数(也就是打印closing websocket的部分)。但你的代码中因为变量名错误导致未捕获异常,打断了这个传递链,同时WebSocket的错误处理缺失,导致流卡住或重试逻辑无法终止。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 18:35:27