本文章分兩大章:
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è)置方可有效。