RxJS concat操作符问题:合并两个Subject仅收到前一个输出
问题分析与解决方案
嘿,我来帮你拆解下这个问题的核心原因,以及对应的解决办法~
为什么只收到1、2、3、4?
这是两个关键点共同作用的结果:
concat的执行逻辑:concat会按顺序依次订阅传入的Observable,只有当前一个Observable完全完成后,才会订阅下一个Observable。在你的代码里,concat(Obs2, Obs1)会先订阅Obs2,直到Obs2调用complete(),才会去订阅Obs1。- 普通
Subject的特性:Subject是一个“热”Observable,它只会把订阅之后发送的值推送给订阅者,订阅前发送的所有值都会直接丢失,不会被缓存。
再理一遍你的代码执行流程:
- 你调用
concat(Obs2, Obs1).subscribe(...)后,concat立刻订阅Obs2,但此时还没订阅Obs1。 - 同步执行
sendToObs1(5)、sendToObs1(6)、sendToObs1(7)——这时候Obs1还没被concat订阅,所以这三个值直接丢失了。 - 10ms后,
sendToObs2完成异步操作,给Obs2发送1、2、3、4并调用Obs2.complete()。 - concat收到Obs2的完成信号后,才去订阅Obs1,但此时Obs1已经没有新值发送了,所以订阅者收不到5、6、7。
解决方案
根据你的业务需求,有两种常见的解决方式:
方案1:用ReplaySubject替代普通Subject
ReplaySubject会缓存指定数量的历史值,当新的订阅者到来时,会把缓存的值重播给订阅者。把Obs1改成ReplaySubject就能保留订阅前发送的5、6、7:
let Obs1 = new rxjs.ReplaySubject(); // 替换为ReplaySubject let Obs2 = new rxjs.Subject(); function sendToObs1(x){ Obs1.next(x) } async function sendToObs2(){ let trns = await getValues(); for(let i = 0; i < trns.length; i++){ Obs2.next(trns[i]) } Obs2.complete() } function getValues(){ return new Promise((resolve, reject) => { setTimeout(() => resolve([1,2,3,4]), 10) }) }; rxjs.concat(Obs2, Obs1).subscribe({ next: x=> console.log("Received: " + x), complete: () => console.log("Done") } ) sendToObs2() sendToObs1(5) sendToObs1(6) sendToObs1(7)
方案2:调整Obs1的发送时机
如果你不需要保留Obs1的历史值,只是想确保Obs1的发送在Obs2完成之后,可以把sendToObs1的调用放在Obs2完成的回调里:
let Obs1 = new rxjs.Subject(); let Obs2 = new rxjs.Subject(); function sendToObs1(x){ Obs1.next(x) } async function sendToObs2(){ let trns = await getValues(); for(let i = 0; i < trns.length; i++){ Obs2.next(trns[i]) } Obs2.complete(); // Obs2完成后再发送Obs1的值 sendToObs1(5); sendToObs1(6); sendToObs1(7); } function getValues(){ return new Promise((resolve, reject) => { setTimeout(() => resolve([1,2,3,4]), 10) }) }; rxjs.concat(Obs2, Obs1).subscribe({ next: x=> console.log("Received: " + x), complete: () => console.log("Done") } ) sendToObs2()
这两种方案都能让你得到预期的输出:1、2、3、4、5、6、7,你可以根据实际业务场景选择合适的方式~
内容的提问来源于stack exchange,提问作者koopatroopa
相关产品推荐
相关产品推荐

