You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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链的逻辑。

在你的代码里:

  1. 第一次订阅s1Observable时,clock发射的0触发map(l->new S1()),生成一个S1实例(比如idS1=1)并打印;
  2. 第二次订阅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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 19:45:17