使用AtomicInteger构造RxJava Observable时的异常行为排查
问题描述
测试代码如下:
import io.reactivex.rxjava3.core.*; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; public class MainTest { public static AtomicInteger s1Id = new AtomicInteger(1); public static AtomicInteger s2Id = new AtomicInteger(1); public static int s1Id() {return s1Id.getAndIncrement();} public static int s2Id() {return s2Id.getAndIncrement();} public static class S1 { public int idS1; public S1() { idS1 = s1Id(); } } public static class S2 { public int idS1; public int idS2; public S2(int idS1) { this.idS1 = idS1; idS2 = s2Id(); } } public static Observable<Long> clock = Observable.interval(1, TimeUnit.SECONDS); public static void main(String args[]) { Observable<S1> s1Observable = clock.filter(l->l<1).map(l->(new S1())); s1Observable.subscribe(s1->System.out.println("S1: "+s1.idS1)); Observable<S2> s2Observable = s1Observable.map(s1->new S2(s1.idS1)); s2Observable.subscribe(s2->System.out.println("S2: idS1: "+s2.idS1 + " idS2: " + s2.idS2)); try{Thread.sleep(100000);}catch(Exception e){e.printStackTrace();} } }
实际输出通常为:
S1: 1
S2: idS1: 2 idS2: 1
偶尔也会出现:
S1: 2
S2: idS1: 1 idS2: 1
按逻辑,S2的idS1应该和对应的S1的idS1一致,但实际不符。引入clock是触发问题的必要条件,请问这里忽略了什么问题?
问题原因与解决办法
核心原因:RxJava冷订阅的特性
RxJava中默认的Observable是冷订阅类型,意思是每调用一次subscribe(),就会从头执行一遍整个Observable链的逻辑。
在你的代码里:
- 第一次订阅
s1Observable时,clock发射的0触发map(l->new S1()),生成一个S1实例(比如idS1=1)并打印; - 第二次订阅
s2Observable时,因为s2Observable基于s1Observable,这会再次触发s1Observable的完整执行链:clock再次发射0,生成另一个新的S1实例(idS1=2),再映射成S2。
这就导致两次订阅分别生成了两个独立的S1实例,自然idS1值不一致——这不是并发问题,是RxJava冷订阅的特性导致的。
解决办法
要让多个订阅共享同一个Observable链的执行结果,需要把冷Observable转换成热Observable,常用方式有两种:
- 使用
publish()+connect():将Observable转为可连接类型,手动触发执行; - 使用
share():是publish().refCount()的简写,会自动管理订阅连接。
修改后的核心代码示例:
// 将s1Observable转为热Observable,共享执行结果 Observable<S1> s1Observable = clock.filter(l->l<1).map(l->(new S1())).share(); s1Observable.subscribe(s1->System.out.println("S1: "+s1.idS1)); Observable<S2> s2Observable = s1Observable.map(s1->new S2(s1.idS1)); s2Observable.subscribe(s2->System.out.println("S2: idS1: "+s2.idS1 + " idS2: " + s2.idS2));
修改后,s1Observable只会执行一次,生成一个S1实例,两个订阅都会拿到这个实例,idS1就能保持一致了。
内容的提问来源于stack exchange,提问作者David Posner
相关产品推荐
相关产品推荐

