RxJS 5.5.6中Observable.empty()配合flatMap未触发订阅回调
我来帮你拆解这个问题,以及为什么会出现和预期不符的行为:
问题根源
首先得明确flatMap(RxJS 5里和mergeMap等价)的工作逻辑:它会把源Observable发射的每个值,转换成一个新的内部Observable,然后订阅这个内部Observable,只把内部Observable的输出(next/error/complete)转发到最终的Observable。
回到你的代码:
- 源Observable是
interval(100).startWith(0):首先发射0,之后每隔100ms持续发射递增数字,这是个无限Observable,永远不会主动触发complete。 - 你在
flatMap里返回了empty():这个Observable的行为是不发射任何next值,直接触发complete通知。
所以当源发射0时,flatMap订阅了empty(),但这个内部Observable不会产生任何next输出,只会告诉flatMap它完成了。接着源Observable会继续发射1、2……每个值都会重复这个过程——订阅新的empty(),没有next输出,只有内部complete。
最终导致:
- 没有
next日志:因为所有内部Observable都不发射任何值; - 没有
complete日志:因为源Observable是无限的,永远不会结束,flatMap会一直等待新的值,整个合并后的Observable也就永远不会触发complete。
而当你换成Rx.Observable.of({})时,这个内部Observable会先发射一个{}的next值,再触发complete,所以flatMap会把这个next转发出去,你的next回调就会被触发(每次源发射值都会触发一次),但complete回调依然不会触发——还是因为源是无限的。
实现预期行为的解决方案
如果你想同时触发next和complete回调(也就是看到0 1和3的日志),需要调整逻辑:
方案1:让源Observable只发射一次就结束
如果不需要后续的interval发射,可以用take(1)限制源只取第一个值(也就是startWith的0),然后让flatMap返回能发射源值的Observable:
var obs = Rx.Observable.interval(100) .startWith(0) .take(1) // 只保留第一个发射值,之后源Observable触发complete .flatMap((val) => { return Rx.Observable.of(val); // 把源的0发射出去,再触发complete }); var sub = obs.subscribe( (data) => {console.log(data, 1)}, (err) => {console.log(err, 2)}, () => {console.log(3)} );
运行后会输出:
0 1 3
方案2:直接去掉flatMap(如果不需要转换逻辑)
如果你的flatMap只是想传递源的值,那完全可以不用flatMap,直接订阅源Observable:
var obs = Rx.Observable.interval(100) .startWith(0) .take(1); var sub = obs.subscribe( (data) => {console.log(data, 1)}, (err) => {console.log(err, 2)}, () => {console.log(3)} );
同样能得到预期的输出。
内容的提问来源于stack exchange,提问作者Saif Adnan

