diff --git a/src/main/java/io/reactivex/rxjava3/core/Observable.java b/src/main/java/io/reactivex/rxjava3/core/Observable.java index c86972a7ad..22f7751eab 100644 --- a/src/main/java/io/reactivex/rxjava3/core/Observable.java +++ b/src/main/java/io/reactivex/rxjava3/core/Observable.java @@ -2643,6 +2643,7 @@ public static Observable intervalRange(long start, long count, long initia @CheckReturnValue @NonNull @SchedulerSupport(SchedulerSupport.NONE) + // 工厂方法,生产一个Observable的子类-ObservableJust public static <@NonNull T> Observable just(@NonNull T item) { Objects.requireNonNull(item, "item is null"); return RxJavaPlugins.onAssembly(new ObservableJust<>(item)); @@ -10373,6 +10374,9 @@ public final Single lastOrError() { @NonNull public final <@NonNull R> Observable map(@NonNull Function mapper) { Objects.requireNonNull(mapper, "mapper is null"); + // 返回一个新的Observable,ObservableMap + // ObservableMap是上游的订阅者 + // ObservableMap是下游的发布者 return RxJavaPlugins.onAssembly(new ObservableMap<>(this, mapper)); } @@ -10611,6 +10615,7 @@ public final Observable observeOn(@NonNull Scheduler scheduler, boolean delay @CheckReturnValue @SchedulerSupport(SchedulerSupport.CUSTOM) @NonNull + // 指定下游的运行线程 public final Observable observeOn(@NonNull Scheduler scheduler, boolean delayError, int bufferSize) { Objects.requireNonNull(scheduler, "scheduler is null"); ObjectHelper.verifyPositive(bufferSize, "bufferSize"); @@ -13166,6 +13171,7 @@ public final Disposable subscribe( @SchedulerSupport(SchedulerSupport.NONE) @Override + // 被Observer订阅 public final void subscribe(@NonNull Observer observer) { Objects.requireNonNull(observer, "observer is null"); try { @@ -13173,6 +13179,7 @@ public final void subscribe(@NonNull Observer observer) { Objects.requireNonNull(observer, "The RxJavaPlugins.onSubscribe hook returned a null Observer. Please change the handler provided to RxJavaPlugins.setOnObservableSubscribe for invalid null returns. Further reading: https://github.com/ReactiveX/RxJava/wiki/Plugins"); + // 模板方法模式,子类实现subscribeActual方法,不需要关注异常 subscribeActual(observer); } catch (NullPointerException e) { // NOPMD throw e; @@ -13196,6 +13203,8 @@ public final void subscribe(@NonNull Observer observer) { * applied by {@link #subscribe(Observer)} before this method gets called. * @param observer the incoming {@code Observer}, never {@code null} */ + // 抽象方法,让子类去实现 + // 模版方法模式 protected abstract void subscribeActual(@NonNull Observer observer); /** @@ -13250,6 +13259,7 @@ public final void subscribe(@NonNull Observer observer) { @CheckReturnValue @SchedulerSupport(SchedulerSupport.CUSTOM) @NonNull + // 切换上游的执行线程 public final Observable subscribeOn(@NonNull Scheduler scheduler) { Objects.requireNonNull(scheduler, "scheduler is null"); return RxJavaPlugins.onAssembly(new ObservableSubscribeOn<>(this, scheduler)); diff --git a/src/main/java/io/reactivex/rxjava3/core/Observer.java b/src/main/java/io/reactivex/rxjava3/core/Observer.java index 6b911f51e5..1ea767fc9c 100644 --- a/src/main/java/io/reactivex/rxjava3/core/Observer.java +++ b/src/main/java/io/reactivex/rxjava3/core/Observer.java @@ -73,6 +73,7 @@ * @param * the type of item the Observer expects to observe */ +// 订阅者的抽象接口 public interface Observer<@NonNull T> { /** diff --git a/src/main/java/io/reactivex/rxjava3/disposables/Disposable.java b/src/main/java/io/reactivex/rxjava3/disposables/Disposable.java index 67fd235086..3f9ba7ace4 100644 --- a/src/main/java/io/reactivex/rxjava3/disposables/Disposable.java +++ b/src/main/java/io/reactivex/rxjava3/disposables/Disposable.java @@ -25,6 +25,8 @@ /** * Represents a disposable resource. */ +// 订阅的句柄 +// 可以取消订阅、可以查询订阅状态(订阅是否被取消) public interface Disposable { /** * Dispose the resource, the operation should be idempotent. diff --git a/src/main/java/io/reactivex/rxjava3/internal/disposables/EmptyDisposable.java b/src/main/java/io/reactivex/rxjava3/internal/disposables/EmptyDisposable.java index fdb30932ed..8a2c5b909c 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/disposables/EmptyDisposable.java +++ b/src/main/java/io/reactivex/rxjava3/internal/disposables/EmptyDisposable.java @@ -25,21 +25,25 @@ * don't use it in tests and then signal onNext with it; * use Disposables.empty() instead. */ +// 单例模式的java枚举写法 public enum EmptyDisposable implements QueueDisposable { /** * Since EmptyDisposable implements QueueDisposable and is empty, * don't use it in tests and then signal onNext with it; * use Disposables.empty() instead. */ + // isDisposed方法永远返回true INSTANCE, /** * An empty disposable that returns false for isDisposed. */ + // isDisposed方法永远返回false NEVER ; @Override public void dispose() { + // dispose的时候什么也不做 // no-op } diff --git a/src/main/java/io/reactivex/rxjava3/internal/fuseable/HasUpstreamObservableSource.java b/src/main/java/io/reactivex/rxjava3/internal/fuseable/HasUpstreamObservableSource.java index 1a1e878b65..43c06df71b 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/fuseable/HasUpstreamObservableSource.java +++ b/src/main/java/io/reactivex/rxjava3/internal/fuseable/HasUpstreamObservableSource.java @@ -22,6 +22,7 @@ * * @param the value type */ +// 有上游ObservableSource的抽象类接口 public interface HasUpstreamObservableSource<@NonNull T> { /** * Returns the upstream source of this Observable. @@ -29,5 +30,6 @@ public interface HasUpstreamObservableSource<@NonNull T> { * @return the source ObservableSource */ @NonNull + // 获取上游的ObservableSource ObservableSource source(); } diff --git a/src/main/java/io/reactivex/rxjava3/internal/observers/BasicFuseableObserver.java b/src/main/java/io/reactivex/rxjava3/internal/observers/BasicFuseableObserver.java index 6bd12f3940..edeb8b92c2 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/observers/BasicFuseableObserver.java +++ b/src/main/java/io/reactivex/rxjava3/internal/observers/BasicFuseableObserver.java @@ -25,6 +25,9 @@ * @param the upstream value type * @param the downstream value type */ +// 中间的observer +// 有一个下游的Observer(downstream) +// 上游的Observable调用它的onSubscribe时,会传递上游的Disposable public abstract class BasicFuseableObserver implements Observer, QueueDisposable { /** The downstream subscriber. */ @@ -53,6 +56,7 @@ public BasicFuseableObserver(Observer downstream) { // final: fixed protocol steps to support fuseable and non-fuseable upstream @SuppressWarnings("unchecked") @Override + // 被上游的Observable调用,传递了上游的Disposable public final void onSubscribe(Disposable d) { if (DisposableHelper.validate(this.upstream, d)) { @@ -63,6 +67,7 @@ public final void onSubscribe(Disposable d) { if (beforeDownstream()) { + // 调用下游的onSubscribe downstream.onSubscribe(this); afterDownstream(); @@ -149,6 +154,7 @@ protected final int transitiveBoundaryFusion(int mode) { @Override public void dispose() { + // 被下游调用dispose后,直接调用上游的dispose upstream.dispose(); } diff --git a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/AbstractObservableWithUpstream.java b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/AbstractObservableWithUpstream.java index dcad29a5d8..d2f0b43187 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/AbstractObservableWithUpstream.java +++ b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/AbstractObservableWithUpstream.java @@ -22,9 +22,11 @@ * @param the input source type * @param the output type */ +// 一个抽象的有上游的Observable abstract class AbstractObservableWithUpstream extends Observable implements HasUpstreamObservableSource { /** The source consumable Observable. */ + // 上游的ObservableSource protected final ObservableSource source; /** diff --git a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableJust.java b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableJust.java index af01eb09bb..a80a67a3ba 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableJust.java +++ b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableJust.java @@ -21,6 +21,9 @@ * Represents a constant scalar value. * @param the value type */ +// Observable的一个子类 +// 被订阅后会自动执行subscribeActual方法 +// 只需要给观察者发送一个value值(调用Observer的onSubscribe方法) public final class ObservableJust extends Observable implements ScalarSupplier { private final T value; @@ -31,6 +34,7 @@ public ObservableJust(final T value) { @Override protected void subscribeActual(Observer observer) { ScalarDisposable sd = new ScalarDisposable<>(observer, value); + // 调用监听者的onSubscribe方法,将Disposable传给监听者 observer.onSubscribe(sd); sd.run(); } diff --git a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableMap.java b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableMap.java index 8f4502ba19..f5bda9d9e3 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableMap.java +++ b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableMap.java @@ -20,9 +20,13 @@ import java.util.Objects; +// 有上游的Observable +// 将上游的T类型的数据,通过function方法,转换成U类型的数据 public final class ObservableMap extends AbstractObservableWithUpstream { final Function function; + // 将上游的ObservableSource包裹起来📦 + // 隔离下游与上游的联系 public ObservableMap(ObservableSource source, Function function) { super(source); this.function = function; @@ -30,6 +34,13 @@ public ObservableMap(ObservableSource source, Function t) { + // 中间层承上启下的作用 + // 下游调用subscribe方法时传递了下游的Observer + // 中间层包裹这个Observer -> MapObserver + // 将MapObserver传递给上游 + // 偷梁换柱,将下游的Observer换成了自己的MapObserver + // 上游调用onNext的时候,MapObserver会将value值进行map转换 + // 转换后的结果再传递给下游的onNext source.subscribe(new MapObserver(t, function)); } @@ -55,11 +66,13 @@ public void onNext(T t) { U v; try { + // 转换value值 v = Objects.requireNonNull(mapper.apply(t), "The mapper function returned a null value."); } catch (Throwable ex) { fail(ex); return; } + // 将转换后的值传递给下游 downstream.onNext(v); } diff --git a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableObserveOn.java b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableObserveOn.java index 05a0064f41..88a32f6ae1 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableObserveOn.java +++ b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableObserveOn.java @@ -43,6 +43,7 @@ protected void subscribeActual(Observer observer) { } else { Scheduler.Worker w = scheduler.createWorker(); + // 将下游的 observer 封装成 ObserveOnObserver,传递给上游 source.subscribe(new ObserveOnObserver<>(observer, w, delayError, bufferSize)); } } @@ -79,6 +80,7 @@ static final class ObserveOnObserver extends BasicIntQueueDisposable @Override public void onSubscribe(Disposable d) { if (DisposableHelper.validate(this.upstream, d)) { + // 保存上游的 Disposable this.upstream = d; if (d instanceof QueueDisposable) { @SuppressWarnings("unchecked") @@ -104,6 +106,7 @@ public void onSubscribe(Disposable d) { queue = new SpscLinkedArrayQueue<>(bufferSize); + // ⚠️下游的 onSubscribe 方法并没有切换线程 downstream.onSubscribe(this); } } @@ -117,6 +120,8 @@ public void onNext(T t) { if (sourceMode != QueueDisposable.ASYNC) { queue.offer(t); } + + // 切换线程,执行下游的 onNext schedule(); } @@ -159,6 +164,7 @@ public boolean isDisposed() { void schedule() { if (getAndIncrement() == 0) { + // 切换线程,执行 run 方法 worker.schedule(this); } } @@ -199,6 +205,7 @@ void drainNormal() { break; } + // 执行下游的 onNext a.onNext(v); } @@ -248,6 +255,7 @@ void drainFused() { } } + // 线程切换以后,执行该方法 @Override public void run() { if (outputFused) { diff --git a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableScalarXMap.java b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableScalarXMap.java index 7ce5bb82f0..dc02aec9bc 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableScalarXMap.java +++ b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableScalarXMap.java @@ -173,6 +173,9 @@ public void subscribeActual(Observer observer) { * * @param the value type */ + // 是一个AtomicInteger,使用CAS保证多线程安全 + // 是一个Runnable,可以被run + // 是一个Disposable,可以被取消执行 public static final class ScalarDisposable extends AtomicInteger implements QueueDisposable, Runnable { @@ -225,11 +228,15 @@ public void clear() { @Override public void dispose() { + // 取消执行 + // 将状态设置为ON_COMPLETE set(ON_COMPLETE); } @Override public boolean isDisposed() { + // 判断是否被取消 + // 判断状态是否是ON_COMPLETE return get() == ON_COMPLETE; } @@ -244,10 +251,13 @@ public int requestFusion(int mode) { @Override public void run() { + // 保证run只会执行一次(使用CAS原子操作) if (get() == START && compareAndSet(START, ON_NEXT)) { + // 执行观察者的onNext方法 observer.onNext(value); if (get() == ON_NEXT) { lazySet(ON_COMPLETE); + // 执行观察者的onComplete方法 observer.onComplete(); } } diff --git a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableSubscribeOn.java b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableSubscribeOn.java index 5eaeaa5035..5ed731633c 100644 --- a/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableSubscribeOn.java +++ b/src/main/java/io/reactivex/rxjava3/internal/operators/observable/ObservableSubscribeOn.java @@ -19,6 +19,7 @@ import io.reactivex.rxjava3.disposables.Disposable; import io.reactivex.rxjava3.internal.disposables.DisposableHelper; +// public final class ObservableSubscribeOn extends AbstractObservableWithUpstream { final Scheduler scheduler; @@ -31,11 +32,23 @@ public ObservableSubscribeOn(ObservableSource source, Scheduler scheduler) { public void subscribeActual(final Observer observer) { final SubscribeOnObserver parent = new SubscribeOnObserver<>(observer); + // 给下游传递当前的Disposable observer.onSubscribe(parent); + // new SubscribeTask(parent)为线程切换后需要执行的task + // scheduleDirect 方法用于切换线程,同时返回一个 Disposable + // 如果切换还没有完成,下游就调用了 dispose 方法,那么就应该取消执行scheduleDirect + // parent.setDisposable 设置当前的 Disposable,给下游调用 parent.setDisposable(scheduler.scheduleDirect(new SubscribeTask(parent))); } + // 由于涉及到线程的切换,在不同的线程,Disposable应该是不同的 + // 所以继承了AtomicReference + // + // SubscribeOnObserver 是一个 Disposable,使下游可以调用 dispose 方法 + // 下游调用 dispose 方法以后, + // 需要 dispose 线程切换 + // 需要 dispose 上游 static final class SubscribeOnObserver extends AtomicReference implements Observer, Disposable { private static final long serialVersionUID = 8094547886072529208L; @@ -50,6 +63,7 @@ static final class SubscribeOnObserver extends AtomicReference im @Override public void onSubscribe(Disposable d) { + // 管理上游的Disposable,如果下游调用了 dispose,则也需要调用上游的 dispose方法 DisposableHelper.setOnce(this.upstream, d); } @@ -70,20 +84,28 @@ public void onComplete() { @Override public void dispose() { + // 下游调用dispose后 + // 先调用上游的dispose + // 在调用当前的dispose DisposableHelper.dispose(upstream); DisposableHelper.dispose(this); } @Override public boolean isDisposed() { + // 只需要判断线程切换有没有被 isDisposed + // 不需要判断上游的 isDisposed return DisposableHelper.isDisposed(get()); } + // 设置线程切换的 Disposable void setDisposable(Disposable d) { DisposableHelper.setOnce(this, d); } } + // 在指定的线程,调用上游的subscribe + // 这样上游的运行线程就被切换了 final class SubscribeTask implements Runnable { private final SubscribeOnObserver parent;