RxJava和RxAndroid是什么?
RxJava是基于JVM的響應式擴展,用于編寫異步代碼
RxAndroid是關于Android的RxJava綁定
RxJava和RxAndroid使用
依賴
implementation 'io.reactivex.rxjava3:rxjava:3.1.0'
implementation 'io.reactivex.rxjava3:rxandroid:3.0.2'
使用過程
如下模擬在子線程中進行耗時操作,并將結果返回到主線程中處理
- Flowable:將要進行的操作
- subscribeOn():操作要運行的線程
- observeOn() :處理結果要運行的線程
- subscribe():處理結果
public class MainActivity extends AppCompatActivity {@Overrideprotected void onCreate(Bundle savedInstanceState) {super.onCreate(savedInstanceState);setContentView(R.layout.activity_main);Flowable.fromCallable(() -> {Thread.sleep(3000);return "Done";}).subscribeOn(Schedulers.io()).observeOn(AndroidSchedulers.mainThread()).subscribe(System.out::println, Throwable::printStackTrace);}
}
內存泄漏問題
若Activity退出后,線程仍未執行完,會導致內存泄漏,需要在onDestroy()將任務取消
public class MainActivity extends AppCompatActivity {CompositeDisposable mCompositeDisposable = new CompositeDisposable();@Overrideprotected void onCreate(Bundle savedInstanceState) {super.onCreate(savedInstanceState);setContentView(R.layout.activity_main);Disposable disposable = Flowable.fromCallable(() -> {Thread.sleep(3000);return "Done";}).subscribeOn(Schedulers.io()).observeOn(AndroidSchedulers.mainThread()).subscribe(System.out::println, Throwable::printStackTrace);mCompositeDisposable.add(disposable);}@Overrideprotected void onDestroy() {super.onDestroy();mCompositeDisposable.dispose();}
}
RxJava源碼解析
Publisher
Publisher用于發布數據,Subscriber通過subscribe()訂閱數據
public interface Publisher<T> {public void subscribe(Subscriber<? super T> s);
}
Subscriber
Subscriber接收Publisher發布的數據
- onSubscribe():subscribe()回調函數,回調前會創建Subscription用于控制數據發布和停止
- onNext():當Subscription調用request()時會調用onNext()發布數據
- onError():處理接收到的錯誤
- onComplete():處理完成的情況
public interface Subscriber<T> {public void onSubscribe(Subscription s);public void onNext(T t);public void onError(Throwable t);public void onComplete();
}
Subscription
Subscription表示Publisher和Subscriber的對應關系
- request():向Publisher請求數據
- cancel():讓Publisher停止發布數據
public interface Subscription {public void request(long n);public void cancel();
}
Scheduler
createWorker()用于創建Worker ,具體的調度工作由Worker的schedule()完成
public abstract class Scheduler {public abstract Worker createWorker();public abstract static class Worker implements Disposable {@NonNullpublic Disposable schedule(@NonNull Runnable run) {return schedule(run, 0L, TimeUnit.NANOSECONDS);}public abstract Disposable schedule(@NonNull Runnable run, long delay, @NonNull TimeUnit unit);}
}
source傳遞過程
fromCallable()創建FlowableFromCallable,傳遞callable
public abstract class Flowable<@NonNull T> implements Publisher<T> {public static <@NonNull T> Flowable<T> fromCallable(@NonNull Callable<? extends T> callable) {return RxJavaPlugins.onAssembly(new FlowableFromCallable<>(callable));}
}
subscribeOn()創建FlowableSubscribeOn,傳遞this(即FlowableFromCallable)作為source
public abstract class Flowable<@NonNull T> implements Publisher<T> {public final Flowable<T> subscribeOn(@NonNull Scheduler scheduler) {Objects.requireNonNull(scheduler, "scheduler is null");return subscribeOn(scheduler, !(this instanceof FlowableCreate));}public final Flowable<T> subscribeOn(@NonNull Scheduler scheduler, boolean requestOn) {Objects.requireNonNull(scheduler, "scheduler is null");return RxJavaPlugins.onAssembly(new FlowableSubscribeOn<>(this, scheduler, requestOn));}
}
observeOn()創建FlowableObserveOn,傳遞this(即FlowableSubscribeOn)作為source
public abstract class Flowable<@NonNull T> implements Publisher<T> {public final Flowable<T> observeOn(@NonNull Scheduler scheduler) {return observeOn(scheduler, false, bufferSize());}public final Flowable<T> observeOn(@NonNull Scheduler scheduler, boolean delayError, int bufferSize) {Objects.requireNonNull(scheduler, "scheduler is null");ObjectHelper.verifyPositive(bufferSize, "bufferSize");return RxJavaPlugins.onAssembly(new FlowableObserveOn<>(this, scheduler, delayError, bufferSize));}
}
即依次將自身當作Flowable,作為參數source傳遞給下一個Flowable
subscribe()流程
subscribe()最終調用具體Flowable的subscribeActual()
public abstract class Flowable<@NonNull T> implements Publisher<T> {......public final Disposable subscribe(@NonNull Consumer<? super T> onNext, @NonNull Consumer<? super Throwable> onError) {return subscribe(onNext, onError, Functions.EMPTY_ACTION);}public final Disposable subscribe(@NonNull Consumer<? super T> onNext, @NonNull Consumer<? super Throwable> onError,@NonNull Action onComplete) {.....LambdaSubscriber<T> ls = new LambdaSubscriber<>(onNext, onError, onComplete, FlowableInternalHelper.RequestMax.INSTANCE);subscribe(ls);return ls;}public final void subscribe(@NonNull FlowableSubscriber<? super T> subscriber) {try {Subscriber<? super T> flowableSubscriber = RxJavaPlugins.onSubscribe(this, subscriber);......subscribeActual(flowableSubscriber);}......}protected abstract void subscribeActual(@NonNull Subscriber<? super T> subscriber);
}
調用過程和傳遞過程是相反的,先調用FlowableObserveOn的subscribeActual()
public final class FlowableObserveOn<T> extends AbstractFlowableWithUpstream<T, T> {......@Overridepublic void subscribeActual(Subscriber<? super T> s) {Worker worker = scheduler.createWorker();if (s instanceof ConditionalSubscriber) {.....} else {source.subscribe(new ObserveOnSubscriber<>(s, worker, delayError, prefetch));}}}
上面的source就是上一層傳遞下來的FlowableSubscribeOn,即調用到FlowableSubscribeOn的subscribeActual()
public final class FlowableSubscribeOn<T> extends AbstractFlowableWithUpstream<T , T> {......@Overridepublic void subscribeActual(final Subscriber<? super T> s) {Scheduler.Worker w = scheduler.createWorker();final SubscribeOnSubscriber<T> sos = new SubscribeOnSubscriber<>(s, w, source, nonScheduledRequests);.....w.schedule(sos);}static final class SubscribeOnSubscriber<T> extends AtomicReference<Thread>implements FlowableSubscriber<T>, Subscription, Runnable {......@Overridepublic void run() {lazySet(Thread.currentThread());Publisher<T> src = source;source = null;src.subscribe(this);}
schedule()最終會調用run()方法,lazySet()切換線程,上面的source就是上一層傳遞下來的FlowableFromCallable,即將到FlowableFromCallable的subscribeActual()放到指定線程中運行
public final class FlowableFromCallable<T> extends Flowable<T> implements Supplier<T> {final Callable<? extends T> callable;......@Overridepublic void subscribeActual(Subscriber<? super T> s) {DeferredScalarSubscription<T> deferred = new DeferredScalarSubscription<>(s);s.onSubscribe(deferred);T t;try {t = Objects.requireNonNull(callable.call(), "The callable returned a null value");} catch (Throwable ex) {Exceptions.throwIfFatal(ex);if (deferred.isCancelled()) {RxJavaPlugins.onError(ex);} else {s.onError(ex);}return;}deferred.complete(t);}......
}
上面若出錯回調onError(),否則調用downstream的onNext()傳遞結果
public class DeferredScalarSubscription<@NonNull T> extends BasicIntQueueSubscription<T> {public final void complete(T v) {int state = get();for (;;) {......if (state == HAS_REQUEST_NO_VALUE) {lazySet(HAS_REQUEST_HAS_VALUE);Subscriber<? super T> a = downstream;a.onNext(v);if (get() != CANCELLED) {a.onComplete();}return;}value = v;......}}
}
onNext()過程
調用FlowableSubscribeOn.SubscribeOnSubscriber的onNext(),調用downstream的onNext()傳遞結果
public final class FlowableSubscribeOn<T> extends AbstractFlowableWithUpstream<T , T> {static final class SubscribeOnSubscriber<T> extends AtomicReference<Thread>implements FlowableSubscriber<T>, Subscription, Runnable {@Overridepublic void onNext(T t) {downstream.onNext(t);}}
}
調用FlowableObserveOn.BaseObserveOnSubscriber的onNext()、trySchedule(),schedule()最終會調用run()方法,根據sourceMode判斷是同步還是異步
- FlowableObserveOn.ObserveOnSubscriber的runSync()和runAsync()都調用downstream的onNext()傳遞結果
public final class FlowableObserveOn<T> extends AbstractFlowableWithUpstream<T, T> {.....abstract static class BaseObserveOnSubscriber<T>extends BasicIntQueueSubscription<T>implements FlowableSubscriber<T>, Runnable {@Overridepublic final void onNext(T t) {......trySchedule();}'final void trySchedule() {......worker.schedule(this);}@Overridepublic final void run() {if (outputFused) {runBackfused();} else if (sourceMode == SYNC) {runSync();} else {runAsync();}}}static final class ObserveOnSubscriber<T> extends BaseObserveOnSubscriber<T>implements FlowableSubscriber<T> {void runSync() {......final Subscriber<? super T> a = downstream;......for (;;) {......while (e != r) {......a.onNext(v);......}......}}.....@Overridevoid runAsync() {final Subscriber<? super T> a = downstream;for (;;) {......while (e != r) {.....a.onNext(v);.....}.....}}}
}
調用LambdaSubscriber的onNext(),通過傳入的Consumer消費掉最終的結果,即通過System.out::println打印出來
public final class LambdaSubscriber<T> extends AtomicReference<Subscription>implements FlowableSubscriber<T>, Subscription, Disposable, LambdaConsumerIntrospection {.....@Overridepublic void onNext(T t) {if (!isDisposed()) {try {onNext.accept(t);} catch (Throwable e) {Exceptions.throwIfFatal(e);get().cancel();onError(e);}}}
}
onSubscribe()和request()流程
FlowableFromCallable回調下一層的onSubscribe(),其將Subscription存到upstream
public final class FlowableFromCallable<T> extends Flowable<T> implements Supplier<T> {final Callable<? extends T> callable;@Overridepublic void subscribeActual(Subscriber<? super T> s) {DeferredScalarSubscription<T> deferred = new DeferredScalarSubscription<>(s);s.onSubscribe(deferred);......}
}public final class FlowableSubscribeOn<T> extends AbstractFlowableWithUpstream<T , T> {static final class SubscribeOnSubscriber<T> extends AtomicReference<Thread>implements FlowableSubscriber<T>, Subscription, Runnable {@Overridepublic void onSubscribe(Subscription s) {if (SubscriptionHelper.setOnce(this.upstream, s)) {......}}}
}
FlowableSubscribeOn回調下一層的onSubscribe(),其回調下一層的onSubscribe()和上一層的request()請求數據
public final class FlowableSubscribeOn<T> extends AbstractFlowableWithUpstream<T , T> {......@Overridepublic void subscribeActual(final Subscriber<? super T> s) {Scheduler.Worker w = scheduler.createWorker();final SubscribeOnSubscriber<T> sos = new SubscribeOnSubscriber<>(s, w, source, nonScheduledRequests);s.onSubscribe(sos);}
}public final class FlowableObserveOn<T> extends AbstractFlowableWithUpstream<T, T> {static final class ObserveOnSubscriber<T> extends BaseObserveOnSubscriber<T>implements FlowableSubscriber<T> {......@Overridepublic void onSubscribe(Subscription s) {if (SubscriptionHelper.validate(this.upstream, s)) {this.upstream = s;......queue = new SpscArrayQueue<>(prefetch);downstream.onSubscribe(this);s.request(prefetch);}}}
}
LambdaSubscriber利用FlowableInternalHelper.RequestMax的accept()調用上一層的request(),從schedule()獲取數據
public final class LambdaSubscriber<T> extends AtomicReference<Subscription>implements FlowableSubscriber<T>, Subscription, Disposable, LambdaConsumerIntrospection {@Overridepublic void onSubscribe(Subscription s) {if (SubscriptionHelper.setOnce(this, s)) {try {onSubscribe.accept(this);} catch (Throwable ex) {......}}}}public final class FlowableInternalHelper {public enum RequestMax implements Consumer<Subscription> {INSTANCE;@Overridepublic void accept(Subscription t) {t.request(Long.MAX_VALUE);}}
}public final class FlowableObserveOn<T> extends AbstractFlowableWithUpstream<T, T> {abstract static class BaseObserveOnSubscriber<T>extends BasicIntQueueSubscription<T>implements FlowableSubscriber<T>, Runnable {@Overridepublic final void request(long n) {if (SubscriptionHelper.validate(n)) {BackpressureHelper.add(requested, n);trySchedule();}}final void trySchedule() {......worker.schedule(this);}}
}
FlowableSubscribeOn.SubscribeOnSubscriber的request()、requestUpstream()判斷當前線程,若未切換線程調用schedule()切換線程調用上一層的request()
public final class FlowableSubscribeOn<T> extends AbstractFlowableWithUpstream<T , T> {static final class SubscribeOnSubscriber<T> extends AtomicReference<Thread>implements FlowableSubscriber<T>, Subscription, Runnable {@Overridepublic void request(final long n) {if (SubscriptionHelper.validate(n)) {Subscription s = this.upstream.get();if (s != null) {requestUpstream(n, s);} else {......}}}}void requestUpstream(final long n, final Subscription s) {if (nonScheduledRequests || Thread.currentThread() == get()) {s.request(n);} else {worker.schedule(new Request(s, n));}}static final class Request implements Runnable {......@Overridepublic void run() {upstream.request(n);}}}
}
DeferredScalarSubscription接收到請求后,將值傳給downstream的onNext()
public class DeferredScalarSubscription<@NonNull T> extends BasicIntQueueSubscription<T> {@Overridepublic final void request(long n) {if (SubscriptionHelper.validate(n)) {for (;;) {int state = get();......if (state == NO_REQUEST_HAS_VALUE) {if (compareAndSet(NO_REQUEST_HAS_VALUE, HAS_REQUEST_HAS_VALUE)) {T v = value;if (v != null) {value = null;Subscriber<? super T> a = downstream;a.onNext(v);if (get() != CANCELLED) {a.onComplete();}}}return;}......}}}
}
Schedulers.io()調度過程
Schedulers.io() = Schedulers.IO = IOTask() = IoHolder.DEFAULT = IoScheduler()
public final class Schedulers {static final Scheduler IO;static final class IoHolder {static final Scheduler DEFAULT = new IoScheduler();}static {IO = RxJavaPlugins.initIoScheduler(new IOTask());}public static Scheduler io() {return RxJavaPlugins.onIoScheduler(IO);}static final class IOTask implements Supplier<Scheduler> {@Overridepublic Scheduler get() {return IoHolder.DEFAULT;}}
}
FlowableSubscribeOn的subscribeActual()通過IoScheduler創建Worker并調用schedule()
public final class FlowableSubscribeOn<T> extends AbstractFlowableWithUpstream<T , T> {......@Overridepublic void subscribeActual(final Subscriber<? super T> s) {Scheduler.Worker w = scheduler.createWorker();final SubscribeOnSubscriber<T> sos = new SubscribeOnSubscriber<>(s, w, source, nonScheduledRequests);.....w.schedule(sos);}
}
調用IoScheduler的createWorker()會返回EventLoopWorker
public final class IoScheduler extends Scheduler {@NonNull@Overridepublic Worker createWorker() {return new EventLoopWorker(pool.get());}
}
調用IoScheduler.EventLoopWorker的schedule()最終調用ThreadWorker的父類NewThreadWorker的scheduleActual()
public final class IoScheduler extends Scheduler {static final class EventLoopWorker extends Scheduler.Worker implements Runnable {@NonNull@Overridepublic Disposable schedule(@NonNull Runnable action, long delayTime, @NonNull TimeUnit unit) {......return threadWorker.scheduleActual(action, delayTime, unit, tasks);}}static final class ThreadWorker extends NewThreadWorker {......}
}
調用scheduleActual()將Runnable封裝成ScheduledRunnable,通過ScheduledThreadPoolExecutor的submit()或schedule()提交
public class NewThreadWorker extends Scheduler.Worker implements Disposable {private final ScheduledExecutorService executor;volatile boolean disposed;public NewThreadWorker(ThreadFactory threadFactory) {executor = SchedulerPoolFactory.create(threadFactory);}@NonNullpublic ScheduledRunnable scheduleActual(final Runnable run, long delayTime, @NonNull TimeUnit unit, @Nullable DisposableContainer parent) {Runnable decoratedRun = RxJavaPlugins.onSchedule(run);ScheduledRunnable sr = new ScheduledRunnable(decoratedRun, parent);......Future<?> f;try {if (delayTime <= 0) {f = executor.submit((Callable<Object>)sr);} else {f = executor.schedule((Callable<Object>)sr, delayTime, unit);}sr.setFuture(f);} catch (RejectedExecutionException ex) {......}return sr;}
}public final class SchedulerPoolFactory {public static ScheduledExecutorService create(ThreadFactory factory) {final ScheduledThreadPoolExecutor exec = new ScheduledThreadPoolExecutor(1, factory);exec.setRemoveOnCancelPolicy(PURGE_ENABLED);return exec;}
}
線程池會調用FlowableSubscribeOn.SubscribeOnSubscriber的run()方法,SubscribeOnSubscriber繼承了AtomicReference<Thread>,lazySet()切換線程調用上一層source的subscribe()
public final class FlowableSubscribeOn<T> extends AbstractFlowableWithUpstream<T , T> {static final class SubscribeOnSubscriber<T> extends AtomicReference<Thread>implements FlowableSubscriber<T>, Subscription, Runnable {@Overridepublic void run() {lazySet(Thread.currentThread());Publisher<T> src = source;source = null;src.subscribe(this);}}}
AndroidSchedulers.mainThread()調度過程
AndroidSchedulers.mainThread() = AndroidSchedulers.MAIN_THREAD = MainHolder.DEFAULT = HandlerScheduler(),通過主線程Looper創建handler
public final class AndroidSchedulers {private static final class MainHolder {static final Scheduler DEFAULT = internalFrom(Looper.getMainLooper(), true);}private static final Scheduler MAIN_THREAD =RxAndroidPlugins.initMainThreadScheduler(() -> MainHolder.DEFAULT);}public static Scheduler mainThread() {return RxAndroidPlugins.onMainThreadScheduler(MAIN_THREAD);}private static Scheduler internalFrom(Looper looper, boolean async) {......return new HandlerScheduler(new Handler(looper), async);}
}
FlowableObserveOn的subscribeActual()通過IoScheduler創建Worker,在onNext()的trySchedule()調用schedule()
public final class FlowableObserveOn<T> extends AbstractFlowableWithUpstream<T, T> {@Overridepublic void subscribeActual(Subscriber<? super T> s) {Worker worker = scheduler.createWorker();if (s instanceof ConditionalSubscriber) {......} else {source.subscribe(new ObserveOnSubscriber<>(s, worker, delayError, prefetch));}}abstract static class BaseObserveOnSubscriber<T>extends BasicIntQueueSubscription<T>implements FlowableSubscriber<T>, Runnable {@Overridepublic final void onNext(T t) {......trySchedule();}final void trySchedule() {......worker.schedule(this);}}
}
調用HandlerScheduler的createWorker()返回HandlerWorker()
final class HandlerScheduler extends Scheduler {@Overridepublic Worker createWorker() {return new HandlerWorker(handler, async);}
}
調用HandlerScheduler.HandlerWorker的schedule(),將Runnable封裝成ScheduledRunnable,調用主線程handler的sendMessageDelayed()
final class HandlerScheduler extends Scheduler {private static final class HandlerWorker extends Worker {@Overridepublic Disposable schedule(Runnable run, long delay, TimeUnit unit) {......run = RxJavaPlugins.onSchedule(run);ScheduledRunnable scheduled = new ScheduledRunnable(handler, run);Message message = Message.obtain(handler, scheduled);message.obj = this; if (async) {message.setAsynchronous(true);}handler.sendMessageDelayed(message, unit.toMillis(delay));......return scheduled;}}
}
最終主線程會調用FlowableObserveOn.BaseObserveOnSubscriber的run(),根據sourceMode判斷是同步還是異步
- FlowableObserveOn.ObserveOnSubscriber的runSync()和runAsync()都調用downstream的onNext()傳遞結果
public final class FlowableObserveOn<T> extends AbstractFlowableWithUpstream<T, T> {.....abstract static class BaseObserveOnSubscriber<T>extends BasicIntQueueSubscription<T>implements FlowableSubscriber<T>, Runnable {......@Overridepublic final void run() {if (outputFused) {runBackfused();} else if (sourceMode == SYNC) {runSync();} else {runAsync();}}}static final class ObserveOnSubscriber<T> extends BaseObserveOnSubscriber<T>implements FlowableSubscriber<T> {void runSync() {......final Subscriber<? super T> a = downstream;......for (;;) {......while (e != r) {......a.onNext(v);......}......}}.....@Overridevoid runAsync() {final Subscriber<? super T> a = downstream;for (;;) {......while (e != r) {.....a.onNext(v);.....}.....}}}
}