flatMap订阅streamToRx生成的Observable永不完成问题排查
你的测试超时的核心问题有两个:一是复用了同一个由streamToRx生成的Observable实例,二是rxjs-stream和你使用的RxJS 4(require('rx')是RxJS 4版本)可能存在兼容性问题,导致文件流结束后,Observable没有正确触发complete通知,进而让测试一直等待done()调用。
下面给你几个可行的解决方案:
方案1:用RxJS 4原生方法转换流(推荐)
RxJS 4本身就内置了Node.js流转Observable的工具Rx.Observable.fromStream,不需要额外依赖rxjs-stream,兼容性拉满:
const Rx = require('rx') const fs = require('fs') it('should not be infinite', done => { Rx.Observable.of(1) // 每次flatMap都创建新的流和Observable,避免复用导致的状态问题 .flatMap(() => Rx.Observable.fromStream(fs.createReadStream('/some/file.txt'))) .map(() => 'file processed') .subscribe( x => console.log('next', x), err => { console.error(err); done(err); }, () => { console.log('complete!'); done(); } ) })
方案2:修复rxjs-stream的使用方式(如果必须保留依赖)
如果你一定要用rxjs-stream,记得把流的创建和Observable转换逻辑放到flatMap内部,确保每次订阅都是全新的流实例(不要在外部提前创建好streamObservable复用):
const Rx = require('rx') const fs = require('fs') const {streamToRx} = require('rxjs-stream') it('should not be infinite', done => { Rx.Observable.of(1) .flatMap(() => streamToRx(fs.createReadStream('/some/file.txt'))) .map(() => 'file processed') .subscribe( x => console.log('next', x), err => { console.error(err); done(err); }, () => { console.log('complete!'); done(); } ) })
额外排查建议
- 先确认
/some/file.txt文件存在且有读取权限,如果文件不存在,流会触发error事件,你可以在错误回调里看到具体报错。 - 如果可以升级RxJS版本,建议转到RxJS 5+,直接用
from操作符就能转换流,完全不需要rxjs-stream:const { of, from } = require('rxjs'); const { flatMap, map } = require('rxjs/operators'); const fs = require('fs'); it('should not be infinite', done => { of(1).pipe( flatMap(() => from(fs.createReadStream('/some/file.txt'))), map(() => 'file processed') ).subscribe( x => console.log('next', x), err => { console.error(err); done(err); }, () => { console.log('complete!'); done(); } ); });
内容的提问来源于stack exchange,提问作者makasprzak
相关产品推荐
相关产品推荐

