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

