基于RxJava-1源碼擴(kuò)展的輪詢器

本文章分兩大章:

1:基于RxJava-1源碼擴(kuò)展的輪詢器

2:NoHttpRxUtils框架與擴(kuò)展輪詢的結(jié)合

基于RxJava-1源碼擴(kuò)展的輪詢器

由于最近NoHttpRxUtils框架使用者反饋NoHttpRxUtils框架怎么沒有輪詢請(qǐng)求的功能呢?正好這幾天有空閑,我就想去實(shí)現(xiàn)這個(gè)功能。前期我想要不要就基于Thread+Runnable+Handler或者Handler+Timer+TimerTask去實(shí)現(xiàn)呢?沉思一會(huì)....就把這兩種方式給否決了。因?yàn)檫@兩種方式實(shí)現(xiàn)的輪詢"對(duì)于我來(lái)說"可控性不高而且產(chǎn)生的代碼邏輯繁瑣。更主要的一點(diǎn)是,NoHttpRxUtils框架是基于RxJava-1的封裝。那么肯定要用RxJava-1的輪詢?nèi)?shí)現(xiàn)。所有我就去看看RxJava-1的輪詢。

RxJava-1的原生輪詢器

    //創(chuàng)建輪詢器
        Observable.interval(3000, 3000, TimeUnit.MILLISECONDS)
                //數(shù)據(jù)處理行動(dòng)監(jiān)聽器--->>此處不可以線程切換
                .map(new Func1<Long, String>() {
                    @Override
                    public String call(Long aLong) {
                        rxJavaNumber++;
                        //在此方法里面做數(shù)據(jù)處理操作。此方法是子線程執(zhí)行的。
                        return "處理完畢的數(shù)據(jù)";
                    }
                })
                //設(shè)置被訂閱者事件執(zhí)行線程-->對(duì).map()"數(shù)據(jù)處理行動(dòng)監(jiān)聽器"線程無(wú)影響
                .subscribeOn(AndroidSchedulers.mainThread())
                //輪詢攔截器
                .takeUntil(new Func1<String, Boolean>() {
                    @Override
                    public Boolean call(String s) {
                        //執(zhí)行10次自動(dòng)停止輪詢
                        if (rxJavaNumber >= 10) {
                            return true;
                        } else {
                            return false;
                        }
                    }
                })
                //訂閱者事件處理器線程切換
                .observeOn(AndroidSchedulers.mainThread())
                //訂閱者事件處理監(jiān)聽器
                .subscribe(new Action1<String>() {
                    @Override
                    public void call(String transferValue) {
                        String toString = mRxjavaPollText.getText().toString();
                        mRxjavaPollText.setText(transferValue + "\n\n" + toString);
                    }
                });

運(yùn)行證實(shí).subscribeOn(AndroidSchedulers.mainThread())的方法設(shè)置無(wú)法對(duì).map()方法里面的實(shí)現(xiàn)類的線程控制,而.map()方法里面實(shí)現(xiàn)類的call()方法卻執(zhí)行在非UI線程中。所以導(dǎo)致RxJava的 Observable.interval(3000, 3000, TimeUnit.MILLISECONDS)無(wú)法像Observable.create(new Observable.OnSubscribe<T>())那樣對(duì)"被觀察者行為監(jiān)聽"和"觀察者事件處理"隨意切換線程,而且在業(yè)務(wù)層次上面好像也有點(diǎn)混亂(以上觀點(diǎn)也許會(huì)由于我對(duì)rxjava不夠深入的了解產(chǎn)生誤差,請(qǐng)大家諒解)。 所以我就對(duì)rxjava-1輪詢?cè)创a進(jìn)行研究(發(fā)現(xiàn)rxjava的線程調(diào)度器真心的牛逼),然后我根據(jù)rxjava輪詢?cè)创a去擴(kuò)展出對(duì)應(yīng)的輪詢

基于RxJava-1源碼擴(kuò)展的輪詢器

rxjava-1輪詢?cè)创a參考類

public final class OnSubscribeTimerPeriodically implements OnSubscribe<Long> {
    final long initialDelay;
    final long period;
    final TimeUnit unit;
    final Scheduler scheduler;

    public OnSubscribeTimerPeriodically(long initialDelay, long period, TimeUnit unit, Scheduler scheduler) {
        this.initialDelay = initialDelay;
        this.period = period;
        this.unit = unit;
        this.scheduler = scheduler;
    }

    @Override
    public void call(final Subscriber<? super Long> child) {
      //由傳入的線程調(diào)度控制Action0實(shí)現(xiàn)類call()方法執(zhí)行在那個(gè)線程中
        final Worker worker = scheduler.createWorker();
        child.add(worker);
        worker.schedulePeriodically(new Action0() {
            long counter;
            @Override
            public void call() {
              //輪詢時(shí),此方法會(huì)觸發(fā)
                try {
                    child.onNext(counter++);
                } catch (Throwable e) {
                    try {
                        worker.unsubscribe();
                    } finally {
                        Exceptions.throwOrReport(e, child);
                    }
                }
            }
            
        }, initialDelay, period, unit);
    }
}

基于參考RxJava-1"OnSubscribeTimerPeriodically"的源碼實(shí)現(xiàn)的輪詢類

public final class OnSubscribeTimerPeriodically<V, T> implements Observable.OnSubscribe<T> {
    /**
     * 初始化加載延遲
     */
    final long initialDelay;
    /**
     * 輪詢間隔時(shí)間
     */
    final long period;
    /**
     * 時(shí)間單位
     */
    final TimeUnit unit;
    /**
     * 訂閱者線程線路
     */
    final Scheduler scheduler;
    /**
     * 可觀察者事件監(jiān)聽器
     */
    private OnObserverEventListener<V, T> eventListener;
    /**
     * 可觀察者線程線路-事件處理默認(rèn)在子線程
     */
    private Scheduler eventScheduler = Schedulers.io();
    /**
     * 傳輸給被觀察者接受的對(duì)象
     */
    private V transferValue;

    public OnSubscribeTimerPeriodically(long initialDelay, long period, TimeUnit unit, Scheduler scheduler) {
        this.initialDelay = initialDelay;
        this.period = period;
        this.unit = unit;
        this.scheduler = scheduler;
    }

    /**
     * 賦值傳輸給被觀察者接受的對(duì)象
     *
     * @param transferValue 被觀察者接受的對(duì)象
     */
    public void setTransferValue(V transferValue) {
        this.transferValue = transferValue;
    }

    /**
     * 設(shè)置事件執(zhí)行線程線路
     *
     * @param eventScheduler 線程線路
     */
    public void setEventScheduler(Scheduler eventScheduler) {
        if (null != eventScheduler) {
            this.eventScheduler = eventScheduler;
        }
    }

    /**
     * 設(shè)置可觀察者事件監(jiān)聽器
     *
     * @param eventListener 可觀察者事件監(jiān)聽器
     */
    public void setOnObserverEventListener(OnObserverEventListener<V, T> eventListener) {
        this.eventListener = eventListener;
    }

    @Override
    public void call(final Subscriber<? super T> subscriber) {
        //子線程調(diào)度器執(zhí)行
        final Scheduler.Worker worker = Schedulers.io().createWorker();
        subscriber.add(worker);
        //開啟子線程調(diào)度器定時(shí)器執(zhí)行
        worker.schedulePeriodically(new Action0() {
            @Override
            public void call() {
                try {
                    synchronized (OnSubscribeTimerPeriodically.this) {
                        //被觀察事件處理的調(diào)度器
                        final Scheduler.Worker workerEvent = eventScheduler.createWorker();
                        workerEvent.schedule(new Action0() {
                            @Override
                            public void call() {
                                if (null != eventListener) {
                                    try {
                                        final T observerEvent = eventListener.onObserverEvent(transferValue);

                                        //觀察者事件處理的調(diào)度器
                                        final Scheduler.Worker workerDispose = scheduler.createWorker();
                                        workerDispose.schedule(new Action0() {
                                            @Override
                                            public void call() {
                                                try {
                                                    subscriber.onNext(observerEvent);
                                                    awake();
                                                } catch (Throwable e) {
                                                    try {
                                                        worker.unsubscribe();
                                                        workerEvent.unsubscribe();
                                                        workerDispose.unsubscribe();
                                                        awake();
                                                    } finally {
                                                        Exceptions.throwOrReport(e, subscriber);
                                                    }
                                                }
                                            }
                                        });
                                    } catch (Throwable e) {
                                        try {
                                            worker.unsubscribe();
                                            workerEvent.unsubscribe();
                                            awake();
                                        } finally {
                                            Exceptions.throwOrReport(e, subscriber);
                                        }
                                    }
                                }
                            }
                        });
                        //當(dāng)前線程休眠,等待"被觀察者"事件邏輯處理完畢
                        OnSubscribeTimerPeriodically.this.wait();
                    }
                } catch (Throwable e) {
                    try {
                        worker.unsubscribe();
                    } finally {
                        Exceptions.throwOrReport(e, subscriber);
                    }
                }
            }

        }, initialDelay, period, unit);
    }

    /**
     * 觀察者已經(jīng)根據(jù)被觀察者的動(dòng)作做出相應(yīng)處理后喚醒調(diào)度器定時(shí)器繼續(xù)往下走
     */
    private void awake() {
        synchronized (OnSubscribeTimerPeriodically.this) {
            OnSubscribeTimerPeriodically.this.notify();
        }
    }
}

OnSubscribeTimerPeriodically執(zhí)行在非UI線程中,能夠通過設(shè)置來(lái)決定"被觀察者"和"觀察者"行為事件執(zhí)行線程

通過繼承RxJava-1"Observable<T>"來(lái)實(shí)現(xiàn)調(diào)用類

public class ObservableExpand<V, T> extends Observable<T> {

    final OnSubscribeTimerPeriodically<V, T> onSubscribeTimerPeriodically;

    /**
     * Creates an Observable with a Function to execute when it is subscribed to.
     * <p>
     * <em>Note:</em> Use {@link #create(OnSubscribe)} to create an Observable, instead of this constructor,
     * unless you specifically have a need for inheritance.
     *
     * @param f {@link OnSubscribe} to be executed when {@link #subscribe(Subscriber)} is called
     */
    protected ObservableExpand(OnSubscribe<T> f, OnSubscribeTimerPeriodically<V, T> onSubscribe1) {
        super(f);
        onSubscribeTimerPeriodically = onSubscribe1;
    }


    /**
     * 輪詢間隔方法
     *
     * @param initialDelay          初始化加載延遲
     * @param period                輪詢間隔時(shí)間
     * @param unit                  時(shí)間單位
     * @param observerEventListener 可觀察者事件監(jiān)聽器
     * @param <T>
     * @return 訂閱間隔計(jì)時(shí)Builder
     */
    public static <V, T> Builder<V, T> intervalPolling(long initialDelay, long period, TimeUnit unit, OnObserverEventListener<V, T> observerEventListener) {

        return new Builder<>(initialDelay, period, unit, observerEventListener);
    }

    /**
     * 訂閱間隔計(jì)時(shí)Builder
     *
     * @param <T>
     */
    public static class Builder<V, T> {
        /**
         * 初始化加載延遲
         */
        private long initialDelay;
        /**
         * 輪詢間隔時(shí)間
         */
        private long period;
        /**
         * 時(shí)間單位
         */
        private TimeUnit unit;
        /**
         * 可觀察者事件監(jiān)聽器
         */
        private OnObserverEventListener<V, T> observerEventListener;

        public Builder(long initialDelay, long period, TimeUnit unit, OnObserverEventListener<V, T> observerEventListener) {
            this.initialDelay = initialDelay;
            this.period = period;
            this.unit = unit;
            this.observerEventListener = observerEventListener;
        }

        /**
         * 設(shè)置可觀察者監(jiān)聽器線程線路
         *
         * @param eventScheduler 線程線路
         * @param transferValue  待處理或者待傳輸?shù)膶?duì)象
         * @return
         */
        public ObservableExpand<V, T> subscribeOn(Scheduler eventScheduler, V transferValue) {

            OnSubscribeTimerPeriodically<V, T> timerPeriodically = new OnSubscribeTimerPeriodically<>(initialDelay, period, unit, Schedulers.computation());
            timerPeriodically.setOnObserverEventListener(observerEventListener);
            timerPeriodically.setTransferValue(transferValue);
            timerPeriodically.setEventScheduler(eventScheduler);
            return new ObservableExpand<>(RxJavaHooks.onCreate(timerPeriodically), timerPeriodically);
        }
    }
}

ObservableExpand繼承與Observable<T>,并采用建筑模式去創(chuàng)建OnSubscribeTimerPeriodically

如何調(diào)用擴(kuò)展的輪詢器?

 ObservableExpand.intervalPolling(3000, 3000, TimeUnit.MILLISECONDS,
                //被觀察者行為監(jiān)聽器 ->正在處理
                new OnObserverEventListener<String, String>() {
                    @Override
                    public String onObserverEvent(String transferValue) {
                        try {
                            //模擬耗時(shí)操作
                            Thread.sleep(3 * 1000);
                            expandNumber++;
                            transferValue = "擴(kuò)展輪詢次數(shù):" + expandNumber + "\n傳輸進(jìn)來(lái)的值:" + transferValue;
                        } catch (InterruptedException e) {
                            e.printStackTrace();
                        }
                        //此處可以放置你要處理的數(shù)據(jù)或者邏輯。
                        return transferValue;
                    }
                })
                //指定被觀察者行為監(jiān)聽器執(zhí)行線程。傳輸對(duì)象進(jìn)被觀察者行為監(jiān)聽器
                .subscribeOn(Schedulers.io(), transitionList)
                //設(shè)置攔截器
                .takeUntil(new Func1<String, Boolean>() {
                    @Override
                    public Boolean call(String untilData) {
                        //執(zhí)行10次自動(dòng)停止輪詢,也可根據(jù)untilData對(duì)象值去判斷是否停止輪詢
                        if (expandNumber >= 10) {
                            return true;
                        } else {
                            return false;
                        }
                    }
                })
                //指定觀察者觸發(fā)監(jiān)聽器執(zhí)行線程
                .observeOn(AndroidSchedulers.mainThread())
                //觀察者觸發(fā)監(jiān)聽器
                .subscribe(new Action1<String>() {
                    @Override
                    public void call(String s) {
                        String toString = mExpandPollText.getText().toString();
                        mExpandPollText.setText(s + "\n\n" + toString);
                    }
                });

擴(kuò)展的輪詢器的調(diào)用就是這么rxjava

點(diǎn)擊進(jìn)入擴(kuò)展輪詢github鏈接

NoHttpRxUtils框架與擴(kuò)展輪詢的結(jié)合

NoHttpRxUtils框架是what?點(diǎn)擊查看NoHttpRxUtils框架博客點(diǎn)擊進(jìn)入NoHttpRxUtils框架github。

NoHttpRxUtils輪詢請(qǐng)求,采用鏈?zhǔn)秸{(diào)用

//獲取請(qǐng)求對(duì)象
RxNoHttpUtils.rxNoHttpRequest()

//此處省略NoHttp網(wǎng)絡(luò)請(qǐng)求設(shè)置參數(shù)的方法
...

//設(shè)置當(dāng)前輪詢請(qǐng)求Sign
.setSign(new Object())

//創(chuàng)建輪詢請(qǐng)求對(duì)象,并指定響應(yīng)轉(zhuǎn)換類型和請(qǐng)求成功或者失敗回調(diào)接口
.builderPoll(Objects.class,new OnIsRequestListener<T>)

//設(shè)置初始化加載延遲
.setInitialDelay(3 * 1000)

//設(shè)置輪詢間隔時(shí)間-默認(rèn)3秒
.setPeriod(5 * 1000)

//設(shè)置被觀察者產(chǎn)生的行為事件監(jiān)聽器-
//(如果此處實(shí)現(xiàn)被觀察者產(chǎn)生的行為事件監(jiān)聽器,那么框架內(nèi)部就不去維護(hù)此輪詢請(qǐng)求,必須實(shí)現(xiàn)輪詢攔截器接口去維護(hù)此輪詢什么時(shí)候停止。)
.setOnObserverEventListener(new OnObserverEventListener<RestRequest<T>, RxInformationModel<T>>(){
     @Override
  public RxInformationModel<T> onObserverEvent(RestRequest<T> transferValue) {
    // RxInformationModel對(duì)象方法介紹
    //getData()=獲取請(qǐng)求數(shù)據(jù)
    //setData(T data)=賦值請(qǐng)求數(shù)據(jù)
    //setException(boolean exception)=賦值是否是異常狀態(tài)
    //isException()=獲取是否異常狀態(tài)
    //setThrowable(Throwable throwable)=賦值異常類
    //getThrowable()=獲取異常類
    //setStop(boolean stop)=賦值是否停止輪詢狀態(tài)
    //isStop()=獲取是否輪詢狀態(tài)

    //RxInformationModel 此對(duì)象需要new 出來(lái).
    //在此方法中可以換成自己鐘意的網(wǎng)絡(luò)框架去請(qǐng)求,如果上面設(shè)置網(wǎng)絡(luò)請(qǐng)求參數(shù),除了body其它的都能從RestRequest里面取得。
   return informationModel;
  }
})

// 設(shè)置設(shè)置數(shù)據(jù)攔截監(jiān)聽對(duì)象
.setBooleanFunc1(new Func1<RxInformationModel<T>, Boolean>() {
    @Override 
   public Boolean call(RxInformationModel<T> stringRxInformationModel) {
  //在此方法里面可以根據(jù)RxInformationModel.getData()獲取請(qǐng)求的數(shù)據(jù),然后根據(jù)請(qǐng)求的數(shù)據(jù)來(lái)決定是否停止輪詢
    return stringRxInformationModel.isStop();
   } })

//設(shè)置觀察者根據(jù)被觀察產(chǎn)生的行為做出相應(yīng)處理監(jiān)聽器
//如果實(shí)現(xiàn)了此接口,那么builderPoll中實(shí)現(xiàn)的OnIsRequestListener將無(wú)效。
.setRxInformationModelAction1(new Action1<RxInformationModel<T>>() {
    @Override 
   public void call(RxInformationModel<T> stringRxInformationModel) {
    //在此方法里面根據(jù)RxInformationModel中的數(shù)據(jù)做出相應(yīng)動(dòng)作
   }
})

//轉(zhuǎn)換成輪詢請(qǐng)求類
.switchPoll()

//開始請(qǐng)求
.requestRxNoHttp();

NoHttpRxUtils輪詢請(qǐng)求調(diào)用中所有的泛型<T>都相互關(guān)聯(lián)的。指定一個(gè)泛型類型,其它泛型都必須是此類型

如何取消輪詢請(qǐng)求

//單個(gè)取消Sign對(duì)應(yīng)的輪詢請(qǐng)求
RxNoHttpUtils.cancelPoll(Sign));

//取消批量Sign對(duì)應(yīng)的輪詢請(qǐng)求
RxNoHttpUtils.cancelPoll(Sign[]);

//取消所有的輪詢請(qǐng)求
// RxNoHttpUtils.cancelPollAll(); 

取消輪詢請(qǐng)求必須要求調(diào)用此方法.setSign(new Object())設(shè)置方可有效。

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
【社區(qū)內(nèi)容提示】社區(qū)部分內(nèi)容疑似由AI輔助生成,瀏覽時(shí)請(qǐng)結(jié)合常識(shí)與多方信息審慎甄別。
平臺(tái)聲明:文章內(nèi)容(如有圖片或視頻亦包括在內(nèi))由作者上傳并發(fā)布,文章內(nèi)容僅代表作者本人觀點(diǎn),簡(jiǎn)書系信息發(fā)布平臺(tái),僅提供信息存儲(chǔ)服務(wù)。

相關(guān)閱讀更多精彩內(nèi)容

友情鏈接更多精彩內(nèi)容