From 560636494efd07d95508e4b83c53a9ef104db73f Mon Sep 17 00:00:00 2001 From: wind0ws <932803995@qq.com> Date: Thu, 12 Jan 2017 16:54:48 +0800 Subject: [PATCH 001/106] Refactor code,more example added. --- app/build.gradle | 29 ++- app/src/main/AndroidManifest.xml | 120 +++++----- .../android/samples/AbsExampleActivity.java | 225 ++++++++++++++++++ .../android/samples/AsyncSubjectExample.java | 131 ---------- .../samples/BehaviorSubjectExample.java | 131 ---------- .../samples/BufferExampleActivity.java | 103 -------- .../samples/CompletableObserverActivity.java | 80 ------- .../samples/ConcatExampleActivity.java | 92 ------- .../samples/DebounceExampleActivity.java | 115 --------- .../android/samples/DeferExampleActivity.java | 90 ------- .../samples/DisposableExampleActivity.java | 98 -------- .../samples/DistinctExampleActivity.java | 75 ------ .../samples/FilterExampleActivity.java | 92 ------- .../samples/FlowableExampleActivity.java | 80 ------- .../samples/IntervalExampleActivity.java | 95 -------- .../samples/LastOperatorExampleActivity.java | 75 ------ .../rxjava2/android/samples/MainActivity.java | 52 +++- .../android/samples/MapExampleActivity.java | 120 ---------- .../android/samples/MergeExampleActivity.java | 90 ------- .../android/samples/MyApplication.java | 30 +++ .../samples/PublishSubjectExample.java | 130 ---------- .../samples/ReduceExampleActivity.java | 90 ------- .../samples/ReplayExampleActivity.java | 129 ---------- .../android/samples/ReplaySubjectExample.java | 129 ---------- .../android/samples/ScanExampleActivity.java | 92 ------- .../samples/SimpleExampleActivity.java | 90 ------- .../SingleObserverExampleActivity.java | 71 ------ .../android/samples/SkipExampleActivity.java | 91 ------- .../android/samples/TakeExampleActivity.java | 91 ------- .../samples/ThrottleLastExampleActivity.java | 119 --------- .../rxjava2/android/samples/TimerExample.java | 92 ------- .../android/samples/ZipExampleActivity.java | 130 ---------- .../android/samples/model/ApiUser.java | 9 + .../rxjava2/android/samples/model/Car.java | 3 +- .../rxjava2/android/samples/model/User.java | 9 + .../AsyncSubjectExampleActivity.java | 39 +++ .../BehaviorSubjectExampleActivity.java | 42 ++++ .../operators/BufferExampleActivity.java | 57 +++++ .../CompletableObserverExampleActivity.java | 47 ++++ .../operators/ConcatExampleActivity.java | 36 +++ .../operators/DebounceExampleActivity.java | 71 ++++++ .../operators/DeferExampleActivity.java | 33 +++ .../operators/DisposableExampleActivity.java | 52 ++++ .../operators/DistinctExampleActivity.java | 20 ++ .../operators/FilterExampleActivity.java | 28 +++ .../operators/FlowableExampleActivity.java | 84 +++++++ .../operators/IntervalExampleActivity.java | 43 ++++ .../LastOperatorExampleActivity.java | 24 ++ .../samples/operators/MapExampleActivity.java | 56 +++++ .../operators/MergeExampleActivity.java | 28 +++ .../PublishSubjectExampleActivity.java | 41 ++++ .../operators/ReduceExampleActivity.java | 31 +++ .../operators/ReplayExampleActivity.java | 45 ++++ .../ReplaySubjectExampleActivity.java | 49 ++++ .../operators/ScanExampleActivity.java | 35 +++ .../operators/SimpleExampleActivity.java | 31 +++ .../SingleObserverExampleActivity.java | 51 ++++ .../operators/SkipExampleActivity.java | 31 +++ .../operators/TakeExampleActivity.java | 32 +++ .../ThrottleFirstExampleActivity.java | 58 +++++ .../ThrottleLastExampleActivity.java | 62 +++++ .../operators/TimerExampleActivity.java | 33 +++ .../operators/WindowExampleActivity.java | 1 + .../samples/operators/ZipExampleActivity.java | 65 +++++ .../rxjava2/android/samples/utils/Utils.java | 13 + app/src/main/res/layout/activity_main.xml | 14 ++ app/src/main/res/values/strings.xml | 2 + build.gradle | 2 +- gradle.properties | 19 +- 69 files changed, 1669 insertions(+), 2804 deletions(-) create mode 100644 app/src/main/java/com/rxjava2/android/samples/AbsExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/AsyncSubjectExample.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/BehaviorSubjectExample.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/BufferExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/CompletableObserverActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/ConcatExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/DebounceExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/DeferExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/DisposableExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/DistinctExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/FilterExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/FlowableExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/IntervalExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/LastOperatorExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/MapExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/MergeExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/MyApplication.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/PublishSubjectExample.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/ReduceExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/ReplayExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/ReplaySubjectExample.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/ScanExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/SimpleExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/SingleObserverExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/SkipExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/TakeExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/ThrottleLastExampleActivity.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/TimerExample.java delete mode 100644 app/src/main/java/com/rxjava2/android/samples/ZipExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/AsyncSubjectExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/BehaviorSubjectExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/BufferExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/CompletableObserverExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/ConcatExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/DebounceExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/DeferExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/DisposableExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/DistinctExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/FilterExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/FlowableExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/IntervalExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/LastOperatorExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/MapExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/MergeExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/PublishSubjectExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/ReduceExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/ReplayExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/ReplaySubjectExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/ScanExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/SimpleExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/SingleObserverExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/SkipExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/TakeExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/ThrottleFirstExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/ThrottleLastExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/TimerExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/WindowExampleActivity.java create mode 100644 app/src/main/java/com/rxjava2/android/samples/operators/ZipExampleActivity.java diff --git a/app/build.gradle b/app/build.gradle index 5955173..5dfd943 100644 --- a/app/build.gradle +++ b/app/build.gradle @@ -1,13 +1,13 @@ apply plugin: 'com.android.application' android { - compileSdkVersion 24 - buildToolsVersion "24.0.1" + compileSdkVersion 25 + buildToolsVersion "25.0.2" defaultConfig { applicationId "com.rxjava2.android.samples" - minSdkVersion 15 - targetSdkVersion 24 + minSdkVersion 16 + targetSdkVersion 25 versionCode 1 versionName "1.0" } @@ -19,10 +19,25 @@ android { } } +//compileOptions { +// sourceCompatibility JavaVersion.VERSION_1_8 +// targetCompatibility JavaVersion.VERSION_1_8 +//} + +//packagingOptions { +// exclude 'LICENSE' +// exclude 'LICENSE.txt' +//} + dependencies { compile fileTree(dir: 'libs', include: ['*.jar']) testCompile 'junit:junit:4.12' - compile 'com.android.support:appcompat-v7:24.1.1' - compile 'io.reactivex.rxjava2:rxjava:2.0.2' - compile 'io.reactivex.rxjava2:rxandroid:2.0.1' + compile 'com.android.support:appcompat-v7:25.1.0' + + // https://mvnrepository.com/artifact/io.reactivex.rxjava2/rxjava + compile group: 'io.reactivex.rxjava2', name: 'rxjava', version: '2.0.4' + // https://mvnrepository.com/artifact/io.reactivex.rxjava2/rxandroid + compile('io.reactivex.rxjava2:rxandroid:2.0.1'){ + exclude group: 'io.reactivex.rxjava2', module: 'rxjava' + } } diff --git a/app/src/main/AndroidManifest.xml b/app/src/main/AndroidManifest.xml index 6f0d4b2..af26764 100644 --- a/app/src/main/AndroidManifest.xml +++ b/app/src/main/AndroidManifest.xml @@ -1,8 +1,9 @@ + package="com.rxjava2.android.samples"> - + - + - + android:name=".operators.SimpleExampleActivity" + android:label="@string/simple"/> + android:name=".operators.MapExampleActivity" + android:label="@string/map"/> + android:name=".operators.ZipExampleActivity" + android:label="@string/zip"/> + android:name=".operators.DisposableExampleActivity" + android:label="@string/disposable"/> + android:name=".operators.TakeExampleActivity" + android:label="@string/take"/> + android:name=".operators.TimerExampleActivity" + android:label="@string/timer"/> + android:name=".operators.IntervalExampleActivity" + android:label="@string/interval"/> + android:name=".operators.SingleObserverExampleActivity" + android:label="@string/SingleObserver"/> + android:name=".operators.CompletableObserverExampleActivity" + android:label="@string/CompletableObserver"/> + android:name=".operators.FlowableExampleActivity" + android:label="@string/Flowable"/> + android:name=".operators.ReduceExampleActivity" + android:label="@string/reduce"/> + android:name=".operators.BufferExampleActivity" + android:label="@string/buffer"/> + android:name=".operators.FilterExampleActivity" + android:label="@string/filter"/> + android:name=".operators.SkipExampleActivity" + android:label="@string/skip"/> + android:name=".operators.ScanExampleActivity" + android:label="@string/scan"/> + android:name=".operators.ReplayExampleActivity" + android:label="@string/replay"/> + android:name=".operators.ConcatExampleActivity" + android:label="@string/concat"/> + android:name=".operators.MergeExampleActivity" + android:label="@string/merge"/> + android:name=".operators.DeferExampleActivity" + android:label="@string/defer"/> + android:name=".operators.DistinctExampleActivity" + android:label="@string/distinct"/> + android:name=".operators.LastOperatorExampleActivity" + android:label="@string/distinct"/> + android:name=".operators.ReplaySubjectExampleActivity" + android:label="@string/replay_subject"/> + android:name=".operators.PublishSubjectExampleActivity" + android:label="@string/publish_subject"/> + android:name=".operators.BehaviorSubjectExampleActivity" + android:label="@string/behavior_subject"/> + android:name=".operators.AsyncSubjectExampleActivity" + android:label="@string/async_subject"/> + android:name=".operators.ThrottleLastExampleActivity" + android:label="@string/throttle_last"/> + android:name=".operators.DebounceExampleActivity" + android:label="@string/debounce"/> + + \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/AbsExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/AbsExampleActivity.java new file mode 100644 index 0000000..d0704d9 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/AbsExampleActivity.java @@ -0,0 +1,225 @@ +package com.rxjava2.android.samples; + +import android.os.Bundle; +import android.support.v7.app.AppCompatActivity; +import android.util.Log; +import android.view.View; +import android.widget.Button; +import android.widget.TextView; + +import com.rxjava2.android.samples.utils.AppConstant; + +import java.util.List; + +import io.reactivex.CompletableObserver; +import io.reactivex.MaybeObserver; +import io.reactivex.Observer; +import io.reactivex.SingleObserver; +import io.reactivex.disposables.Disposable; +import io.reactivex.observers.DisposableObserver; + +/** + * Created by threshold on 2017/1/11. + */ + +public abstract class AbsExampleActivity extends AppCompatActivity { + + protected String TAG = getClass().getSimpleName(); + protected Button btn; + protected TextView textView; + + @Override + protected void onCreate(Bundle savedInstanceState) { + super.onCreate(savedInstanceState); + setContentView(R.layout.activity_example); + btn = (Button) findViewById(R.id.btn); + textView = (TextView) findViewById(R.id.textView); + btn.setOnClickListener(new View.OnClickListener() { + @Override + public void onClick(View view) { + doSomeWork(); + } + }); + } + + /** + * Do some work when click the button. + */ + protected abstract void doSomeWork() ; + + protected Observer getObserver() { + return getObserver(""); + } + + protected Observer getObserver(final Object theNumberOfObserver) { + return new Observer() { + + @Override + public void onSubscribe(Disposable d) { + AbsExampleActivity.this.onSubscribe(theNumberOfObserver,d); + } + + @Override + public void onNext(T value) { + AbsExampleActivity.this.onNext(theNumberOfObserver,value); + } + + @Override + public void onError(Throwable e) { + AbsExampleActivity.this.onError(theNumberOfObserver,e); + } + + @Override + public void onComplete() { + AbsExampleActivity.this.onComplete(theNumberOfObserver); + } + }; + } + + protected CompletableObserver getCompletableObserver() { + return getCompletableObserver(""); + } + + protected CompletableObserver getCompletableObserver(final Object theNumberOfObserver) { + return new CompletableObserver() { + @Override + public void onSubscribe(Disposable d) { + AbsExampleActivity.this.onSubscribe(theNumberOfObserver, d); + } + + @Override + public void onComplete() { + AbsExampleActivity.this.onComplete(theNumberOfObserver); + } + + @Override + public void onError(Throwable e) { + AbsExampleActivity.this.onError(theNumberOfObserver, e); + } + }; + } + + protected DisposableObserver getDisposableObserver() { + return getDisposableObserver(""); + } + + protected DisposableObserver getDisposableObserver(final Object theNumberOfObserver) { + return new DisposableObserver() { + @Override + public void onNext(T t) { + AbsExampleActivity.this.onNext(theNumberOfObserver, t); + } + + @Override + public void onError(Throwable e) { + AbsExampleActivity.this.onError(theNumberOfObserver, e); + } + + @Override + public void onComplete() { + AbsExampleActivity.this.onComplete(theNumberOfObserver); + } + }; + } + + protected SingleObserver getSingleObserver() { + return getSingleObserver(""); + } + + protected SingleObserver getSingleObserver(final Object theNumberOfObserver) { + return new SingleObserver() { + @Override + public void onSubscribe(Disposable d) { + AbsExampleActivity.this.onSubscribe(theNumberOfObserver, d); + } + + @Override + public void onSuccess(T t) { + AbsExampleActivity.this.onNext(theNumberOfObserver, t); + } + + @Override + public void onError(Throwable e) { + AbsExampleActivity.this.onError(theNumberOfObserver, e); + } + }; + } + + protected MaybeObserver getMaybeObserver() { + return getMaybeObserver(""); + } + + protected MaybeObserver getMaybeObserver(final Object theNumberOfObserver) { + return new MaybeObserver() { + @Override + public void onSubscribe(Disposable d) { + AbsExampleActivity.this.onSubscribe(theNumberOfObserver, d); + } + + @Override + public void onSuccess(T t) { + AbsExampleActivity.this.onSuccess(theNumberOfObserver, t); + } + + @Override + public void onError(Throwable e) { + AbsExampleActivity.this.onError(theNumberOfObserver, e); + } + + @Override + public void onComplete() { + AbsExampleActivity.this.onComplete(theNumberOfObserver); + } + }; + } + + protected void onSubscribe(Object theNumberOfObserver,Disposable d) { + String msg = theNumberOfObserver + " onSubscribe"; + if (d != null) { + msg = msg + ": isDisposed :" + d.isDisposed(); + } + Log.d(TAG, msg); + textView.append(msg); + textView.append(AppConstant.LINE_SEPARATOR); + } + + protected void onNext(Object theNumberOfObserver,T value) { + onPositive(theNumberOfObserver,"onNext",value); + } + + protected void onSuccess(Object theNumberOfObserver, T value) { + onPositive(theNumberOfObserver,"onSuccess",value); + } + + private void onPositive(Object theNumberOfObserver,String positiveKey,T value) { + if (value instanceof List) { + List valueList = (List) value; + Log.d(TAG, theNumberOfObserver+" "+positiveKey+" size :" + valueList.size()); + textView.append( theNumberOfObserver+" "+positiveKey+" size :" + valueList.size()); + textView.append(AppConstant.LINE_SEPARATOR); + for (Object obj : valueList) { + Log.d(TAG, " : value : " + obj.toString()); + textView.append(" value : " + obj.toString()); + textView.append(AppConstant.LINE_SEPARATOR); + } + } else { + Log.d(TAG, theNumberOfObserver+" "+positiveKey+" value : " + value); + textView.append(theNumberOfObserver+" "+positiveKey+" : value : " + value); + textView.append(AppConstant.LINE_SEPARATOR); + } + } + + protected void onError(Object theNumberOfObserver,Throwable e) { + textView.append(theNumberOfObserver+" onError : " + e.getMessage()); + textView.append(AppConstant.LINE_SEPARATOR); + Log.e(TAG, theNumberOfObserver+" onError : " , e); + } + + public void onComplete(Object theNumberOfObserver) { + textView.append(theNumberOfObserver+" onComplete"); + textView.append(AppConstant.LINE_SEPARATOR); + Log.d(TAG, theNumberOfObserver+" onComplete"); + } + + +} diff --git a/app/src/main/java/com/rxjava2/android/samples/AsyncSubjectExample.java b/app/src/main/java/com/rxjava2/android/samples/AsyncSubjectExample.java deleted file mode 100644 index 0b15b73..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/AsyncSubjectExample.java +++ /dev/null @@ -1,131 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; -import io.reactivex.subjects.AsyncSubject; - -/** - * Created by amitshekhar on 17/12/16. - */ - -public class AsyncSubjectExample extends AppCompatActivity { - - private static final String TAG = AsyncSubjectExample.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* An AsyncSubject emits the last value (and only the last value) emitted by the source - * Observable, and only after that source Observable completes. (If the source Observable - * does not emit any values, the AsyncSubject also completes without emitting any values.) - */ - private void doSomeWork() { - - AsyncSubject source = AsyncSubject.create(); - - source.subscribe(getFirstObserver()); // it will emit only 4 and onComplete - - source.onNext(1); - source.onNext(2); - source.onNext(3); - - /* - * it will emit 4 and onComplete for second observer also. - */ - source.subscribe(getSecondObserver()); - - source.onNext(4); - source.onComplete(); - - } - - - private Observer getFirstObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " First onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" First onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" First onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" First onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onComplete"); - } - }; - } - - private Observer getSecondObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - textView.append(" Second onSubscribe : isDisposed :" + d.isDisposed()); - Log.d(TAG, " Second onSubscribe : " + d.isDisposed()); - textView.append(AppConstant.LINE_SEPARATOR); - } - - @Override - public void onNext(Integer value) { - textView.append(" Second onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" Second onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" Second onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/BehaviorSubjectExample.java b/app/src/main/java/com/rxjava2/android/samples/BehaviorSubjectExample.java deleted file mode 100644 index 8644d2f..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/BehaviorSubjectExample.java +++ /dev/null @@ -1,131 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; -import io.reactivex.subjects.BehaviorSubject; - -/** - * Created by amitshekhar on 17/12/16. - */ - -public class BehaviorSubjectExample extends AppCompatActivity { - - private static final String TAG = BehaviorSubjectExample.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* When an observer subscribes to a BehaviorSubject, it begins by emitting the item most - * recently emitted by the source Observable (or a seed/default value if none has yet been - * emitted) and then continues to emit any other items emitted later by the source Observable(s). - */ - private void doSomeWork() { - - BehaviorSubject source = BehaviorSubject.create(); - - source.subscribe(getFirstObserver()); // it will get 1, 2, 3, 4 and onComplete - - source.onNext(1); - source.onNext(2); - source.onNext(3); - - /* - * it will emit 3(last emitted), 4 and onComplete for second observer also. - */ - source.subscribe(getSecondObserver()); - - source.onNext(4); - source.onComplete(); - - } - - - private Observer getFirstObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " First onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" First onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" First onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" First onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onComplete"); - } - }; - } - - private Observer getSecondObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - textView.append(" Second onSubscribe : isDisposed :" + d.isDisposed()); - Log.d(TAG, " Second onSubscribe : " + d.isDisposed()); - textView.append(AppConstant.LINE_SEPARATOR); - } - - @Override - public void onNext(Integer value) { - textView.append(" Second onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" Second onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" Second onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/BufferExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/BufferExampleActivity.java deleted file mode 100644 index 8e2c9a7..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/BufferExampleActivity.java +++ /dev/null @@ -1,103 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import java.util.List; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class BufferExampleActivity extends AppCompatActivity { - - private static final String TAG = BufferExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * simple example using buffer operator - bundles all emitted values into a list - */ - private void doSomeWork() { - - Observable> buffered = getObservable().buffer(3, 1); - - // 3 means, it takes max of three from its start index and create list - // 1 means, it jumps one step every time - // so the it gives the following list - // 1 - one, two, three - // 2 - two, three, four - // 3 - three, four, five - // 4 - four, five - // 5 - five - - buffered.subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.just("one", "two", "three", "four", "five"); - } - - private Observer> getObserver() { - return new Observer>() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(List stringList) { - textView.append(" onNext size : " + stringList.size()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext : size :" + stringList.size()); - for (String value : stringList) { - textView.append(" value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " : value :" + value); - } - - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/CompletableObserverActivity.java b/app/src/main/java/com/rxjava2/android/samples/CompletableObserverActivity.java deleted file mode 100644 index b825a7b..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/CompletableObserverActivity.java +++ /dev/null @@ -1,80 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import java.util.concurrent.TimeUnit; - -import io.reactivex.Completable; -import io.reactivex.CompletableObserver; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class CompletableObserverActivity extends AppCompatActivity { - - private static final String TAG = CompletableObserverActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * simple example using CompletableObserver - */ - private void doSomeWork() { - Completable completable = Completable.timer(1000, TimeUnit.MILLISECONDS); - - completable - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .subscribe(getCompletableObserver()); - } - - private CompletableObserver getCompletableObserver() { - return new CompletableObserver() { - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - }; - } - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/ConcatExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/ConcatExampleActivity.java deleted file mode 100644 index 1af58c5..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/ConcatExampleActivity.java +++ /dev/null @@ -1,92 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class ConcatExampleActivity extends AppCompatActivity { - - private static final String TAG = ConcatExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * Using concat operator to combine Observable : concat maintain - * the order of Observable. - * It will emit all the 7 values in order - * here - first "A1", "A2", "A3", "A4" and then "B1", "B2", "B3" - * first all from the first Observable and then - * all from the second Observable all in order - */ - private void doSomeWork() { - final String[] aStrings = {"A1", "A2", "A3", "A4"}; - final String[] bStrings = {"B1", "B2", "B3"}; - - final Observable aObservable = Observable.fromArray(aStrings); - final Observable bObservable = Observable.fromArray(bStrings); - - Observable.concat(aObservable, bObservable) - .subscribe(getObserver()); - } - - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(String value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext : value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/DebounceExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/DebounceExampleActivity.java deleted file mode 100644 index 411d028..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/DebounceExampleActivity.java +++ /dev/null @@ -1,115 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import java.util.concurrent.TimeUnit; - -import io.reactivex.Observable; -import io.reactivex.ObservableEmitter; -import io.reactivex.ObservableOnSubscribe; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 22/12/16. - */ - -public class DebounceExampleActivity extends AppCompatActivity { - - private static final String TAG = DebounceExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * Using debounce() -> only emit an item from an Observable if a particular time-span has - * passed without it emitting another item, so it will emit 2, 4, 5 as we have simulated it. - */ - private void doSomeWork() { - getObservable() - .debounce(500, TimeUnit.MILLISECONDS) - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.create(new ObservableOnSubscribe() { - @Override - public void subscribe(ObservableEmitter emitter) throws Exception { - // send events with simulated time wait - emitter.onNext(1); // skip - Thread.sleep(400); - emitter.onNext(2); // deliver - Thread.sleep(505); - emitter.onNext(3); // skip - Thread.sleep(100); - emitter.onNext(4); // deliver - Thread.sleep(605); - emitter.onNext(5); // deliver - Thread.sleep(510); - emitter.onComplete(); - } - }); - } - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" onNext : "); - textView.append(AppConstant.LINE_SEPARATOR); - textView.append(" value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext "); - Log.d(TAG, " value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/DeferExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/DeferExampleActivity.java deleted file mode 100644 index 6c0937b..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/DeferExampleActivity.java +++ /dev/null @@ -1,90 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.model.Car; -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; - -/** - * Created by amitshekhar on 30/08/16. - */ -public class DeferExampleActivity extends AppCompatActivity { - - private static final String TAG = DeferExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * Defer used for Deferring Observable code until subscription in RxJava - */ - private void doSomeWork() { - - Car car = new Car(); - - Observable brandDeferObservable = car.brandDeferObservable(); - - car.setBrand("BMW"); // Even if we are setting the brand after creating Observable - // we will get the brand as BMW. - // If we had not used defer, we would have got null as the brand. - - brandDeferObservable - .subscribe(getObserver()); - } - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(String value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext : value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/DisposableExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/DisposableExampleActivity.java deleted file mode 100644 index b581d54..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/DisposableExampleActivity.java +++ /dev/null @@ -1,98 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.os.SystemClock; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import java.util.concurrent.Callable; - -import io.reactivex.Observable; -import io.reactivex.ObservableSource; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.CompositeDisposable; -import io.reactivex.observers.DisposableObserver; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class DisposableExampleActivity extends AppCompatActivity { - - private static final String TAG = DisposableExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - private final CompositeDisposable disposables = new CompositeDisposable(); - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - @Override - protected void onDestroy() { - super.onDestroy(); - disposables.clear(); // do not send event after activity has been destroyed - } - - /* - * Example to understand how to use disposables. - * disposables is cleared in onDestroy of this activity. - */ - void doSomeWork() { - disposables.add(sampleObservable() - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .subscribeWith(new DisposableObserver() { - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onNext(String value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext value : " + value); - } - })); - } - - static Observable sampleObservable() { - return Observable.defer(new Callable>() { - @Override - public ObservableSource call() throws Exception { - // Do some long running operation - SystemClock.sleep(2000); - return Observable.just("one", "two", "three", "four", "five"); - } - }); - } -} - diff --git a/app/src/main/java/com/rxjava2/android/samples/DistinctExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/DistinctExampleActivity.java deleted file mode 100644 index dc8f5e2..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/DistinctExampleActivity.java +++ /dev/null @@ -1,75 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.annotation.Nullable; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; -import com.rxjava2.android.samples.utils.AppConstant; -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; - -/** - * Created by techteam on 13/09/16. - */ -public class DistinctExampleActivity extends AppCompatActivity { - - private static final String TAG = DistinctExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override protected void onCreate(@Nullable Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - - private void doSomeWork(){ - - getObservable().distinct() .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.just(1, 2, 1, 1, 2, 3, 4 ,6, 4); - } - - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - Log.d(TAG, " onComplete"); - } - }; - } -} diff --git a/app/src/main/java/com/rxjava2/android/samples/FilterExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/FilterExampleActivity.java deleted file mode 100644 index 0638d8f..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/FilterExampleActivity.java +++ /dev/null @@ -1,92 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; -import io.reactivex.functions.Predicate; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class FilterExampleActivity extends AppCompatActivity { - - private static final String TAG = FilterExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * simple example by using filter operator to emit only even value - * - */ - private void doSomeWork() { - Observable.just(1, 2, 3, 4, 5, 6) - .filter(new Predicate() { - @Override - public boolean test(Integer integer) throws Exception { - return integer % 2 == 0; - } - }) - .subscribe(getObserver()); - } - - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" onNext : "); - textView.append(AppConstant.LINE_SEPARATOR); - textView.append(" value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext "); - Log.d(TAG, " value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/FlowableExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/FlowableExampleActivity.java deleted file mode 100644 index cbead49..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/FlowableExampleActivity.java +++ /dev/null @@ -1,80 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Flowable; -import io.reactivex.SingleObserver; -import io.reactivex.disposables.Disposable; -import io.reactivex.functions.BiFunction; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class FlowableExampleActivity extends AppCompatActivity { - - private static final String TAG = FlowableExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * simple example using Flowable - */ - private void doSomeWork() { - - Flowable observable = Flowable.just(1, 2, 3, 4); - - observable.reduce(50, new BiFunction() { - @Override - public Integer apply(Integer t1, Integer t2) { - return t1 + t2; - } - }).subscribe(getObserver()); - - } - - private SingleObserver getObserver() { - - return new SingleObserver() { - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onSuccess(Integer value) { - textView.append(" onSuccess : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onSuccess : value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - }; - } -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/IntervalExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/IntervalExampleActivity.java deleted file mode 100644 index afaf822..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/IntervalExampleActivity.java +++ /dev/null @@ -1,95 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import java.util.concurrent.TimeUnit; - -import io.reactivex.Observable; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.CompositeDisposable; -import io.reactivex.observers.DisposableObserver; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class IntervalExampleActivity extends AppCompatActivity { - - private static final String TAG = IntervalExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - private final CompositeDisposable disposables = new CompositeDisposable(); - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - @Override - protected void onDestroy() { - super.onDestroy(); - disposables.clear(); // clearing it : do not emit after destroy - } - - /* - * simple example using interval to run task at an interval of 2 sec - * which start immediately - */ - private void doSomeWork() { - disposables.add(getObservable() - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .subscribeWith(getObserver())); - } - - private Observable getObservable() { - return Observable.interval(0, 2, TimeUnit.SECONDS); - } - - private DisposableObserver getObserver() { - return new DisposableObserver() { - - @Override - public void onNext(Long value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext : value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/LastOperatorExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/LastOperatorExampleActivity.java deleted file mode 100644 index d17a22b..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/LastOperatorExampleActivity.java +++ /dev/null @@ -1,75 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.annotation.Nullable; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.SingleObserver; -import io.reactivex.disposables.Disposable; - -/** - * Created by techteam on 13/09/16. - */ -public class LastOperatorExampleActivity extends AppCompatActivity { - - private static final String TAG = DistinctExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - - @Override - protected void onCreate(@Nullable Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - private void doSomeWork() { - getObservable().last("A1") // the default item ("A1") to emit if the source ObservableSource is empty - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.just("A1", "A2", "A3", "A4", "A5", "A6"); - } - - private SingleObserver getObserver() { - return new SingleObserver() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onSuccess(String value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext value : " + value); - } - - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - }; - } -} diff --git a/app/src/main/java/com/rxjava2/android/samples/MainActivity.java b/app/src/main/java/com/rxjava2/android/samples/MainActivity.java index 9080a31..50ba052 100644 --- a/app/src/main/java/com/rxjava2/android/samples/MainActivity.java +++ b/app/src/main/java/com/rxjava2/android/samples/MainActivity.java @@ -5,6 +5,36 @@ import android.support.v7.app.AppCompatActivity; import android.view.View; +import com.rxjava2.android.samples.operators.AsyncSubjectExampleActivity; +import com.rxjava2.android.samples.operators.BehaviorSubjectExampleActivity; +import com.rxjava2.android.samples.operators.BufferExampleActivity; +import com.rxjava2.android.samples.operators.CompletableObserverExampleActivity; +import com.rxjava2.android.samples.operators.ConcatExampleActivity; +import com.rxjava2.android.samples.operators.DebounceExampleActivity; +import com.rxjava2.android.samples.operators.DeferExampleActivity; +import com.rxjava2.android.samples.operators.DisposableExampleActivity; +import com.rxjava2.android.samples.operators.DistinctExampleActivity; +import com.rxjava2.android.samples.operators.FilterExampleActivity; +import com.rxjava2.android.samples.operators.FlowableExampleActivity; +import com.rxjava2.android.samples.operators.IntervalExampleActivity; +import com.rxjava2.android.samples.operators.LastOperatorExampleActivity; +import com.rxjava2.android.samples.operators.MapExampleActivity; +import com.rxjava2.android.samples.operators.MergeExampleActivity; +import com.rxjava2.android.samples.operators.PublishSubjectExampleActivity; +import com.rxjava2.android.samples.operators.ReduceExampleActivity; +import com.rxjava2.android.samples.operators.ReplayExampleActivity; +import com.rxjava2.android.samples.operators.ReplaySubjectExampleActivity; +import com.rxjava2.android.samples.operators.ScanExampleActivity; +import com.rxjava2.android.samples.operators.SimpleExampleActivity; +import com.rxjava2.android.samples.operators.SingleObserverExampleActivity; +import com.rxjava2.android.samples.operators.SkipExampleActivity; +import com.rxjava2.android.samples.operators.TakeExampleActivity; +import com.rxjava2.android.samples.operators.ThrottleFirstExampleActivity; +import com.rxjava2.android.samples.operators.ThrottleLastExampleActivity; +import com.rxjava2.android.samples.operators.TimerExampleActivity; +import com.rxjava2.android.samples.operators.WindowExampleActivity; +import com.rxjava2.android.samples.operators.ZipExampleActivity; + public class MainActivity extends AppCompatActivity { @Override @@ -34,7 +64,7 @@ public void startTakeActivity(View view) { } public void startTimerActivity(View view) { - startActivity(new Intent(MainActivity.this, TimerExample.class)); + startActivity(new Intent(MainActivity.this, TimerExampleActivity.class)); } public void startIntervalActivity(View view) { @@ -46,7 +76,7 @@ public void startSingleObserverActivity(View view) { } public void startCompletableObserverActivity(View view) { - startActivity(new Intent(MainActivity.this, CompletableObserverActivity.class)); + startActivity(new Intent(MainActivity.this, CompletableObserverExampleActivity.class)); } public void startFlowableActivity(View view) { @@ -98,19 +128,23 @@ public void startLastOperatorActivity(View view) { } public void startReplaySubjectActivity(View view) { - startActivity(new Intent(MainActivity.this, ReplaySubjectExample.class)); + startActivity(new Intent(MainActivity.this, ReplaySubjectExampleActivity.class)); } public void startPublishSubjectActivity(View view) { - startActivity(new Intent(MainActivity.this, PublishSubjectExample.class)); + startActivity(new Intent(MainActivity.this, PublishSubjectExampleActivity.class)); } public void startBehaviorSubjectActivity(View view) { - startActivity(new Intent(MainActivity.this, BehaviorSubjectExample.class)); + startActivity(new Intent(MainActivity.this, BehaviorSubjectExampleActivity.class)); } public void startAsyncSubjectActivity(View view) { - startActivity(new Intent(MainActivity.this, AsyncSubjectExample.class)); + startActivity(new Intent(MainActivity.this, AsyncSubjectExampleActivity.class)); + } + + public void startThrottleFirstActivity(View view) { + startActivity(new Intent(MainActivity.this,ThrottleFirstExampleActivity.class)); } public void startThrottleLastActivity(View view) { @@ -120,4 +154,10 @@ public void startThrottleLastActivity(View view) { public void startDebounceActivity(View view) { startActivity(new Intent(MainActivity.this, DebounceExampleActivity.class)); } + + public void startWindowActivity(View view) { + startActivity(new Intent(MainActivity.this,WindowExampleActivity.class)); + } + + } diff --git a/app/src/main/java/com/rxjava2/android/samples/MapExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/MapExampleActivity.java deleted file mode 100644 index 6ddd151..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/MapExampleActivity.java +++ /dev/null @@ -1,120 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.model.ApiUser; -import com.rxjava2.android.samples.model.User; -import com.rxjava2.android.samples.utils.AppConstant; -import com.rxjava2.android.samples.utils.Utils; - -import java.util.List; - -import io.reactivex.Observable; -import io.reactivex.ObservableEmitter; -import io.reactivex.ObservableOnSubscribe; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.functions.Function; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class MapExampleActivity extends AppCompatActivity { - - private static final String TAG = MapExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * Here we are getting ApiUser Object from api server - * then we are converting it into User Object because - * may be our database support User Not ApiUser Object - * Here we are using Map Operator to do that - */ - private void doSomeWork() { - getObservable() - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .map(new Function, List>() { - - @Override - public List apply(List apiUsers) throws Exception { - return Utils.convertApiUserListToUserList(apiUsers); - } - }) - .subscribe(getObserver()); - } - - private Observable> getObservable() { - return Observable.create(new ObservableOnSubscribe>() { - @Override - public void subscribe(ObservableEmitter> e) throws Exception { - if (!e.isDisposed()) { - e.onNext(Utils.getApiUserList()); - e.onComplete(); - } - } - }); - } - - private Observer> getObserver() { - return new Observer>() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(List userList) { - textView.append(" onNext"); - textView.append(AppConstant.LINE_SEPARATOR); - for (User user : userList) { - textView.append(" firstName : " + user.firstName); - textView.append(AppConstant.LINE_SEPARATOR); - } - Log.d(TAG, " onNext : " + userList.size()); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/MergeExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/MergeExampleActivity.java deleted file mode 100644 index d8ae8c0..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/MergeExampleActivity.java +++ /dev/null @@ -1,90 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; - -/** - * Created by amitshekhar on 28/08/16. - */ -public class MergeExampleActivity extends AppCompatActivity { - - private static final String TAG = MergeExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * Using merge operator to combine Observable : merge does not maintain - * the order of Observable. - * It will emit all the 7 values may not be in order - * Ex - "A1", "B1", "A2", "A3", "A4", "B2", "B3" - may be anything - */ - private void doSomeWork() { - final String[] aStrings = {"A1", "A2", "A3", "A4"}; - final String[] bStrings = {"B1", "B2", "B3"}; - - final Observable aObservable = Observable.fromArray(aStrings); - final Observable bObservable = Observable.fromArray(bStrings); - - Observable.merge(aObservable, bObservable) - .subscribe(getObserver()); - } - - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(String value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext : value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/MyApplication.java b/app/src/main/java/com/rxjava2/android/samples/MyApplication.java new file mode 100644 index 0000000..c2aac36 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/MyApplication.java @@ -0,0 +1,30 @@ +package com.rxjava2.android.samples; + +import android.app.Application; +import android.os.Build; +import android.util.Log; + +import io.reactivex.Observable; +import io.reactivex.functions.Consumer; + +/** + * Created by threshold on 2017/1/12. + */ + +public class MyApplication extends Application { + + @Override + public void onCreate() { + super.onCreate(); + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.LOLLIPOP) { + Log.d("MyApplication","isArm64:" + (Build.SUPPORTED_64_BIT_ABIS.length > 0)); + Observable.fromArray(Build.SUPPORTED_64_BIT_ABIS) + .forEach(new Consumer() { + @Override + public void accept(String s) throws Exception { + Log.d("MyApplication", s); + } + }); + } + } +} diff --git a/app/src/main/java/com/rxjava2/android/samples/PublishSubjectExample.java b/app/src/main/java/com/rxjava2/android/samples/PublishSubjectExample.java deleted file mode 100644 index 2e4b40c..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/PublishSubjectExample.java +++ /dev/null @@ -1,130 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; -import io.reactivex.subjects.PublishSubject; - -/** - * Created by amitshekhar on 17/12/16. - */ - -public class PublishSubjectExample extends AppCompatActivity { - - private static final String TAG = PublishSubjectExample.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* PublishSubject emits to an observer only those items that are emitted - * by the source Observable, subsequent to the time of the subscription. - */ - private void doSomeWork() { - - PublishSubject source = PublishSubject.create(); - - source.subscribe(getFirstObserver()); // it will get 1, 2, 3, 4 and onComplete - - source.onNext(1); - source.onNext(2); - source.onNext(3); - - /* - * it will emit 4 and onComplete for second observer also. - */ - source.subscribe(getSecondObserver()); - - source.onNext(4); - source.onComplete(); - - } - - - private Observer getFirstObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " First onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" First onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" First onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" First onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onComplete"); - } - }; - } - - private Observer getSecondObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - textView.append(" Second onSubscribe : isDisposed :" + d.isDisposed()); - Log.d(TAG, " Second onSubscribe : " + d.isDisposed()); - textView.append(AppConstant.LINE_SEPARATOR); - } - - @Override - public void onNext(Integer value) { - textView.append(" Second onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" Second onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" Second onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/ReduceExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/ReduceExampleActivity.java deleted file mode 100644 index 68fbfd6..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/ReduceExampleActivity.java +++ /dev/null @@ -1,90 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.MaybeObserver; -import io.reactivex.Observable; -import io.reactivex.disposables.Disposable; -import io.reactivex.functions.BiFunction; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class ReduceExampleActivity extends AppCompatActivity { - - private static final String TAG = ReduceExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * simple example using reduce to add all the number - */ - private void doSomeWork() { - getObservable() - .reduce(new BiFunction() { - @Override - public Integer apply(Integer t1, Integer t2) { - return t1 + t2; - } - }) - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.just(1, 2, 3, 4); - } - - private MaybeObserver getObserver() { - return new MaybeObserver() { - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onSuccess(Integer value) { - textView.append(" onSuccess : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onSuccess : value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/ReplayExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/ReplayExampleActivity.java deleted file mode 100644 index 4ce591a..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/ReplayExampleActivity.java +++ /dev/null @@ -1,129 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; -import io.reactivex.observables.ConnectableObservable; -import io.reactivex.subjects.PublishSubject; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class ReplayExampleActivity extends AppCompatActivity { - - private static final String TAG = ReplayExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* Using replay operator, replay ensure that all observers see the same sequence - * of emitted items, even if they subscribe after the Observable has begun emitting items - */ - private void doSomeWork() { - - PublishSubject source = PublishSubject.create(); - ConnectableObservable connectableObservable = source.replay(3); // bufferSize = 3 to retain 3 values to replay - connectableObservable.connect(); // connecting the connectableObservable - - connectableObservable.subscribe(getFirstObserver()); - - source.onNext(1); - source.onNext(2); - source.onNext(3); - source.onNext(4); - source.onComplete(); - - /* - * it will emit 2, 3, 4 as (count = 3), retains the 3 values for replay - */ - connectableObservable.subscribe(getSecondObserver()); - - } - - - private Observer getFirstObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " First onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" First onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" First onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" First onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onComplete"); - } - }; - } - - private Observer getSecondObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - textView.append(" Second onSubscribe : isDisposed :" + d.isDisposed()); - Log.d(TAG, " Second onSubscribe : " + d.isDisposed()); - textView.append(AppConstant.LINE_SEPARATOR); - } - - @Override - public void onNext(Integer value) { - textView.append(" Second onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" Second onError : " + e.getMessage()); - Log.d(TAG, " Second onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" Second onComplete"); - Log.d(TAG, " Second onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/ReplaySubjectExample.java b/app/src/main/java/com/rxjava2/android/samples/ReplaySubjectExample.java deleted file mode 100644 index baedec7..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/ReplaySubjectExample.java +++ /dev/null @@ -1,129 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observer; -import io.reactivex.disposables.Disposable; -import io.reactivex.subjects.ReplaySubject; - -/** - * Created by amitshekhar on 17/12/16. - */ - -public class ReplaySubjectExample extends AppCompatActivity { - - private static final String TAG = ReplaySubjectExample.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* ReplaySubject emits to any observer all of the items that were emitted - * by the source Observable, regardless of when the observer subscribes. - */ - private void doSomeWork() { - - ReplaySubject source = ReplaySubject.create(); - - source.subscribe(getFirstObserver()); // it will get 1, 2, 3, 4 - - source.onNext(1); - source.onNext(2); - source.onNext(3); - source.onNext(4); - source.onComplete(); - - /* - * it will emit 1, 2, 3, 4 for second observer also as we have used replay - */ - source.subscribe(getSecondObserver()); - - } - - - private Observer getFirstObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " First onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" First onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" First onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" First onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " First onComplete"); - } - }; - } - - private Observer getSecondObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - textView.append(" Second onSubscribe : isDisposed :" + d.isDisposed()); - Log.d(TAG, " Second onSubscribe : " + d.isDisposed()); - textView.append(AppConstant.LINE_SEPARATOR); - } - - @Override - public void onNext(Integer value) { - textView.append(" Second onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" Second onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" Second onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " Second onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/ScanExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/ScanExampleActivity.java deleted file mode 100644 index c3ce226..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/ScanExampleActivity.java +++ /dev/null @@ -1,92 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.functions.BiFunction; -import io.reactivex.schedulers.Schedulers; - -public class ScanExampleActivity extends AppCompatActivity { - - private static final String TAG = ScanExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* Using scan operator, it sends also the previous result */ - private void doSomeWork() { - getObservable() - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .scan(new BiFunction() { - @Override - public Integer apply(Integer int1, Integer int2) throws Exception { - return int1 + int2; - } - }) - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.just(1, 2, 3, 4, 5); - } - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/SimpleExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/SimpleExampleActivity.java deleted file mode 100644 index 35fe203..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/SimpleExampleActivity.java +++ /dev/null @@ -1,90 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class SimpleExampleActivity extends AppCompatActivity { - - private static final String TAG = SimpleExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * simple example to emit two value one by one - */ - private void doSomeWork() { - getObservable() - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.just("Cricket", "Football"); - } - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(String value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext : value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/SingleObserverExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/SingleObserverExampleActivity.java deleted file mode 100644 index 68d4aae..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/SingleObserverExampleActivity.java +++ /dev/null @@ -1,71 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Single; -import io.reactivex.SingleObserver; -import io.reactivex.disposables.Disposable; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class SingleObserverExampleActivity extends AppCompatActivity { - - private static final String TAG = SingleObserverExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * simple example using SingleObserver - */ - private void doSomeWork() { - Single.just("Amit") - .subscribe(getSingleObserver()); - } - - private SingleObserver getSingleObserver() { - return new SingleObserver() { - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onSuccess(String value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - }; - } - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/SkipExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/SkipExampleActivity.java deleted file mode 100644 index e12d510..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/SkipExampleActivity.java +++ /dev/null @@ -1,91 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class SkipExampleActivity extends AppCompatActivity { - - private static final String TAG = SkipExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* Using skip operator, it only not emit - * the first 2 values. - */ - private void doSomeWork() { - getObservable() - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .skip(2) - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.just(1, 2, 3, 4, 5); - } - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/TakeExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/TakeExampleActivity.java deleted file mode 100644 index ea55e33..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/TakeExampleActivity.java +++ /dev/null @@ -1,91 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class TakeExampleActivity extends AppCompatActivity { - - private static final String TAG = TakeExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* Using take operator, it only emits - * required number of values. here only 3 out of 5 - */ - private void doSomeWork() { - getObservable() - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .take(3) - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.just(1, 2, 3, 4, 5); - } - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/ThrottleLastExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/ThrottleLastExampleActivity.java deleted file mode 100644 index 2fe83ac..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/ThrottleLastExampleActivity.java +++ /dev/null @@ -1,119 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import java.util.concurrent.TimeUnit; - -import io.reactivex.Observable; -import io.reactivex.ObservableEmitter; -import io.reactivex.ObservableOnSubscribe; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 22/12/16. - */ - -public class ThrottleLastExampleActivity extends AppCompatActivity { - - private static final String TAG = ThrottleLastExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * Using throttleLast() -> emit the most recent items emitted by an Observable within - * periodic time intervals, so here it will emit 2, 6 and 7 as we have simulated it to be the - * last the element in the interval of 500 millis - */ - private void doSomeWork() { - getObservable() - .throttleLast(500, TimeUnit.MILLISECONDS) - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.create(new ObservableOnSubscribe() { - @Override - public void subscribe(ObservableEmitter emitter) throws Exception { - // send events with simulated time wait - Thread.sleep(0); - emitter.onNext(1); // skip - emitter.onNext(2); // deliver - Thread.sleep(505); - emitter.onNext(3); // skip - Thread.sleep(99); - emitter.onNext(4); // skip - Thread.sleep(100); - emitter.onNext(5); // skip - emitter.onNext(6); // deliver - Thread.sleep(305); - emitter.onNext(7); // deliver - Thread.sleep(510); - emitter.onComplete(); - } - }); - } - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Integer value) { - textView.append(" onNext : "); - textView.append(AppConstant.LINE_SEPARATOR); - textView.append(" value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext "); - Log.d(TAG, " value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/TimerExample.java b/app/src/main/java/com/rxjava2/android/samples/TimerExample.java deleted file mode 100644 index 7277f69..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/TimerExample.java +++ /dev/null @@ -1,92 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.utils.AppConstant; - -import java.util.concurrent.TimeUnit; - -import io.reactivex.Observable; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class TimerExample extends AppCompatActivity { - - private static final String TAG = TimerExample.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * simple example using timer to do something after 2 second - */ - private void doSomeWork() { - getObservable() - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .subscribe(getObserver()); - } - - private Observable getObservable() { - return Observable.timer(2, TimeUnit.SECONDS); - } - - private Observer getObserver() { - return new Observer() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(Long value) { - textView.append(" onNext : value : " + value); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onNext : value : " + value); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/ZipExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/ZipExampleActivity.java deleted file mode 100644 index abbcc3d..0000000 --- a/app/src/main/java/com/rxjava2/android/samples/ZipExampleActivity.java +++ /dev/null @@ -1,130 +0,0 @@ -package com.rxjava2.android.samples; - -import android.os.Bundle; -import android.support.v7.app.AppCompatActivity; -import android.util.Log; -import android.view.View; -import android.widget.Button; -import android.widget.TextView; - -import com.rxjava2.android.samples.model.User; -import com.rxjava2.android.samples.utils.AppConstant; -import com.rxjava2.android.samples.utils.Utils; - -import java.util.List; - -import io.reactivex.Observable; -import io.reactivex.ObservableEmitter; -import io.reactivex.ObservableOnSubscribe; -import io.reactivex.Observer; -import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.disposables.Disposable; -import io.reactivex.functions.BiFunction; -import io.reactivex.schedulers.Schedulers; - -/** - * Created by amitshekhar on 27/08/16. - */ -public class ZipExampleActivity extends AppCompatActivity { - - private static final String TAG = ZipExampleActivity.class.getSimpleName(); - Button btn; - TextView textView; - - @Override - protected void onCreate(Bundle savedInstanceState) { - super.onCreate(savedInstanceState); - setContentView(R.layout.activity_example); - btn = (Button) findViewById(R.id.btn); - textView = (TextView) findViewById(R.id.textView); - - btn.setOnClickListener(new View.OnClickListener() { - @Override - public void onClick(View view) { - doSomeWork(); - } - }); - } - - /* - * Here we are getting two user list - * One, the list of cricket fans - * Another one, the list of football fans - * Then we are finding the list of users who loves both - */ - private void doSomeWork() { - Observable.zip(getCricketFansObservable(), getFootballFansObservable(), - new BiFunction, List, List>() { - @Override - public List apply(List cricketFans, List footballFans) throws Exception { - return Utils.filterUserWhoLovesBoth(cricketFans, footballFans); - } - }) - // Run on a background thread - .subscribeOn(Schedulers.io()) - // Be notified on the main thread - .observeOn(AndroidSchedulers.mainThread()) - .subscribe(getObserver()); - } - - private Observable> getCricketFansObservable() { - return Observable.create(new ObservableOnSubscribe>() { - @Override - public void subscribe(ObservableEmitter> e) throws Exception { - if (!e.isDisposed()) { - e.onNext(Utils.getUserListWhoLovesCricket()); - e.onComplete(); - } - } - }); - } - - private Observable> getFootballFansObservable() { - return Observable.create(new ObservableOnSubscribe>() { - @Override - public void subscribe(ObservableEmitter> e) throws Exception { - if (!e.isDisposed()) { - e.onNext(Utils.getUserListWhoLovesFootball()); - e.onComplete(); - } - } - }); - } - - private Observer> getObserver() { - return new Observer>() { - - @Override - public void onSubscribe(Disposable d) { - Log.d(TAG, " onSubscribe : " + d.isDisposed()); - } - - @Override - public void onNext(List userList) { - textView.append(" onNext"); - textView.append(AppConstant.LINE_SEPARATOR); - for (User user : userList) { - textView.append(" firstName : " + user.firstName); - textView.append(AppConstant.LINE_SEPARATOR); - } - Log.d(TAG, " onNext : " + userList.size()); - } - - @Override - public void onError(Throwable e) { - textView.append(" onError : " + e.getMessage()); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onError : " + e.getMessage()); - } - - @Override - public void onComplete() { - textView.append(" onComplete"); - textView.append(AppConstant.LINE_SEPARATOR); - Log.d(TAG, " onComplete"); - } - }; - } - - -} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/model/ApiUser.java b/app/src/main/java/com/rxjava2/android/samples/model/ApiUser.java index 6335019..941bfc3 100644 --- a/app/src/main/java/com/rxjava2/android/samples/model/ApiUser.java +++ b/app/src/main/java/com/rxjava2/android/samples/model/ApiUser.java @@ -7,4 +7,13 @@ public class ApiUser { public long id; public String firstName; public String lastName; + + @Override + public String toString() { + return "ApiUser{" + + "id=" + id + + ", firstName='" + firstName + '\'' + + ", lastName='" + lastName + '\'' + + '}'; + } } diff --git a/app/src/main/java/com/rxjava2/android/samples/model/Car.java b/app/src/main/java/com/rxjava2/android/samples/model/Car.java index e223650..de771e4 100644 --- a/app/src/main/java/com/rxjava2/android/samples/model/Car.java +++ b/app/src/main/java/com/rxjava2/android/samples/model/Car.java @@ -10,13 +10,14 @@ */ public class Car { - private String brand; + private String brand="Volvo"; public void setBrand(String brand) { this.brand = brand; } public Observable brandDeferObservable() { +// return Observable.just(brand); return Observable.defer(new Callable>() { @Override public ObservableSource call() throws Exception { diff --git a/app/src/main/java/com/rxjava2/android/samples/model/User.java b/app/src/main/java/com/rxjava2/android/samples/model/User.java index de5a522..5d64ffb 100644 --- a/app/src/main/java/com/rxjava2/android/samples/model/User.java +++ b/app/src/main/java/com/rxjava2/android/samples/model/User.java @@ -7,4 +7,13 @@ public class User { public long id; public String firstName; public String lastName; + + @Override + public String toString() { + return "User{" + + "id=" + id + + ", firstName='" + firstName + '\'' + + ", lastName='" + lastName + '\'' + + '}'; + } } diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/AsyncSubjectExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/AsyncSubjectExampleActivity.java new file mode 100644 index 0000000..f6f6f47 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/AsyncSubjectExampleActivity.java @@ -0,0 +1,39 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.subjects.AsyncSubject; + +/** + * Created by amitshekhar on 17/12/16. + */ + +public class AsyncSubjectExampleActivity extends AbsExampleActivity { + + /* An AsyncSubject emits the last value (and only the last value) emitted by the source + * Observable, and only after that source Observable completes. (If the source Observable + * does not emit any values, the AsyncSubject also completes without emitting any values.) + */ + protected void doSomeWork() { + AsyncSubject source = AsyncSubject.create(); + + source.onNext(0); + source.subscribe(getObserver("First")); // it will emit only 4 and onComplete + + source.onNext(1); + source.onNext(2); + source.onNext(3); + + /* + * it will emit 4 and onComplete for second observer also. + */ + source.subscribe(getObserver("Second")); + + source.onNext(4); + source.onComplete(); + + source.subscribe(getObserver("Third")); + } + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/BehaviorSubjectExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/BehaviorSubjectExampleActivity.java new file mode 100644 index 0000000..ba86805 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/BehaviorSubjectExampleActivity.java @@ -0,0 +1,42 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.subjects.BehaviorSubject; + +/** + * Created by amitshekhar on 17/12/16. + */ + +public class BehaviorSubjectExampleActivity extends AbsExampleActivity { + + /* When an observer subscribes to a BehaviorSubject, it begins by emitting the item most + * recently emitted by the source Observable (or a seed/default value if none has yet been + * emitted) and then continues to emit any other items emitted later by the source Observable(s). + */ + protected void doSomeWork() { + + BehaviorSubject source = BehaviorSubject.createDefault(-1); + +// source.onNext(0); //如果启用这句话,firstObserver将获得 0 (最近发射的最后一个数据),1,2,3,4 + + source.subscribe(getObserver("First")); // it will get -1, 1, 2, 3, 4 and onComplete + + source.onNext(1); + source.onNext(2); + source.onNext(3); + + /* + * it will emit 3(last emitted), 4 and onComplete for second observer also. + */ + source.subscribe(getObserver("Second")); + + source.onNext(4); + source.onComplete(); + + source.subscribe(getObserver("Third")); + + } + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/BufferExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/BufferExampleActivity.java new file mode 100644 index 0000000..e65479e --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/BufferExampleActivity.java @@ -0,0 +1,57 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class BufferExampleActivity extends AbsExampleActivity { + + /* + * simple example using buffer operator - bundles all emitted values into a list + */ + protected void doSomeWork() { + + /* + + Observable> buffered = getObservable().buffer(3, 1); + + // 3 means, it takes max of three from its start index and create list + // 1 means, it jumps one step every time + // so the it gives the following list + // 1 - one, two, three + // 2 - two, three, four + // 3 - three, four, five + // 4 - four, five + // 5 - five + + buffered.subscribe(getObserver()); + + */ + + /* + //每次取2个 每次跳过3个 + //第一次:one、two + //第二次:four、five + getObservable().buffer(2, 3) + .subscribe(getObserver()); + */ + + //每次取3个,每次跳过2个 + //第一次:one、two、three + //第二次:three、four、five + //第三次:five + getObservable().buffer(3, 2) + .subscribe(getObserver()); + + } + + private Observable getObservable() { + return Observable.just("one", "two", "three", "four", "five"); + } + + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/CompletableObserverExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/CompletableObserverExampleActivity.java new file mode 100644 index 0000000..6a03e0d --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/CompletableObserverExampleActivity.java @@ -0,0 +1,47 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import java.util.Random; + +import io.reactivex.Completable; +import io.reactivex.CompletableEmitter; +import io.reactivex.CompletableOnSubscribe; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class CompletableObserverExampleActivity extends AbsExampleActivity { + + /* + * simple example using CompletableObserver + */ + protected void doSomeWork() { + //延迟多久发射onComplete +// Completable completable = Completable.timer(1000, TimeUnit.MILLISECONDS); + + Completable completable = Completable.create(new CompletableOnSubscribe() { + @Override + public void subscribe(CompletableEmitter e) throws Exception { + if (!e.isDisposed()) { + int randomInt = new Random().nextInt(10); + if (randomInt % 2 == 0) { + e.onComplete(); + } else { + e.onError(new IllegalStateException("Can't completable because of error occur.")); + } + } + } + }); + + completable + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribe(getCompletableObserver()); + } + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/ConcatExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/ConcatExampleActivity.java new file mode 100644 index 0000000..da35280 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/ConcatExampleActivity.java @@ -0,0 +1,36 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class ConcatExampleActivity extends AbsExampleActivity { + + private static final String TAG = ConcatExampleActivity.class.getSimpleName(); + + + /* + * Using concat operator to combine Observable : concat maintain + * the order of Observable. + * It will emit all the 7 values in order + * here - first "A1", "A2", "A3", "A4" and then "B1", "B2", "B3" + * first all from the first Observable and then + * all from the second Observable all in order + */ + protected void doSomeWork() { + final String[] aStrings = {"A1", "A2", "A3", "A4"}; + final String[] bStrings = {"B1", "B2", "B3"}; + + final Observable aObservable = Observable.fromArray(aStrings); + final Observable bObservable = Observable.fromArray(bStrings); + + Observable.concat(aObservable, bObservable) + .subscribe(getObserver()); + } + + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/DebounceExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/DebounceExampleActivity.java new file mode 100644 index 0000000..3ea8656 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/DebounceExampleActivity.java @@ -0,0 +1,71 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import java.util.concurrent.TimeUnit; + +import io.reactivex.Observable; +import io.reactivex.ObservableEmitter; +import io.reactivex.ObservableOnSubscribe; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 22/12/16. + */ + +public class DebounceExampleActivity extends AbsExampleActivity { + + private static final String TAG = DebounceExampleActivity.class.getSimpleName(); + + + /* + * Using debounce() -> only emit an item from an Observable if a particular time-span has + * passed without it emitting another item, so it will emit 2, 4, 5 as we have simulated it. + */ + protected void doSomeWork() { + + //debounce 是发射所有 时间片段 交集 中最后一个元素。 + //这里的时间片段是指每次发射元素的当前时刻+超时时间形成的时间片段 + getObservable() + .debounce(500, TimeUnit.MILLISECONDS) + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribe(getObserver()); + } + + private Observable getObservable() { + return Observable.create(new ObservableOnSubscribe() { + @Override + public void subscribe(ObservableEmitter emitter) throws Exception { + // send events with simulated time wait + emitter.onNext(1); // skip + Thread.sleep(400); + emitter.onNext(2); // deliver + Thread.sleep(505); + emitter.onNext(3); // skip + Thread.sleep(100); + emitter.onNext(4); // deliver + Thread.sleep(605); + emitter.onNext(5); // deliver + Thread.sleep(510); + emitter.onComplete(); + + /* + 分析如下: + 1和2在400毫秒处有交集,所以1被扔掉。 + 2和3之间有500毫秒的间隔,没有交集,所以2被发射出去。 + 3在905毫秒出准备发射,但是紧接着4在1005毫秒处也要准备发射,所以3和4有交集,3被扔掉。 + 4和5之间有605毫秒的间隔,没有交集,所以4被发射出去。 + 5在接下来的500毫秒内没有和其他元素有交集,所以发射出去。 + (如还不明白,建议在纸上画出各个元素的时间片段) + */ + + } + }); + } + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/DeferExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/DeferExampleActivity.java new file mode 100644 index 0000000..dbdf465 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/DeferExampleActivity.java @@ -0,0 +1,33 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; +import com.rxjava2.android.samples.model.Car; + +import io.reactivex.Observable; + +/** + * Created by amitshekhar on 30/08/16. + */ +public class DeferExampleActivity extends AbsExampleActivity { + + /* + * Defer used for Deferring Observable code until subscription in RxJava + */ + protected void doSomeWork() { + + Car car = new Car(); + +// car.setBrand("Audi"); + + Observable brandDeferObservable = car.brandDeferObservable(); + + car.setBrand("BMW"); // Even if we are setting the brand after creating Observable + // we will get the brand as BMW. + // If we had not used defer, we would have got null as the brand. + + brandDeferObservable + .subscribe(getObserver()); + } + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/DisposableExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/DisposableExampleActivity.java new file mode 100644 index 0000000..2846f8f --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/DisposableExampleActivity.java @@ -0,0 +1,52 @@ +package com.rxjava2.android.samples.operators; + +import android.os.SystemClock; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import java.util.concurrent.Callable; + +import io.reactivex.Observable; +import io.reactivex.ObservableSource; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.disposables.CompositeDisposable; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class DisposableExampleActivity extends AbsExampleActivity { + + private final CompositeDisposable disposables = new CompositeDisposable(); + + @Override + protected void onDestroy() { + super.onDestroy(); + disposables.clear(); // do not send event after activity has been destroyed + } + + /* + * Example to understand how to use disposables. + * disposables is cleared in onDestroy of this activity. + */ + protected void doSomeWork() { + disposables.add(sampleObservable() + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribeWith(this.getDisposableObserver())); + } + + static Observable sampleObservable() { + return Observable.defer(new Callable>() { + @Override + public ObservableSource call() throws Exception { + // Do some long running operation + SystemClock.sleep(2000); + return Observable.just("one", "two", "three", "four", "five"); + } + }); + } +} + diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/DistinctExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/DistinctExampleActivity.java new file mode 100644 index 0000000..62f70bd --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/DistinctExampleActivity.java @@ -0,0 +1,20 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; + +/** + * Created by techteam on 13/09/16. + */ +public class DistinctExampleActivity extends AbsExampleActivity { + + protected void doSomeWork(){ + getObservable().distinct() .subscribe(getObserver()); + } + + private Observable getObservable() { + return Observable.just(1, 2, 1, 1, 2, 3, 4 ,6, 4); + } + +} diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/FilterExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/FilterExampleActivity.java new file mode 100644 index 0000000..dc85561 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/FilterExampleActivity.java @@ -0,0 +1,28 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; +import io.reactivex.functions.Predicate; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class FilterExampleActivity extends AbsExampleActivity { + + /* + * simple example by using filter operator to emit only even value + * + */ + protected void doSomeWork() { + Observable.just(1, 2, 3, 4, 5, 6) + .filter(new Predicate() { + @Override + public boolean test(Integer integer) throws Exception { + return integer % 2 == 0; + } + }) + .subscribe(getObserver()); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/FlowableExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/FlowableExampleActivity.java new file mode 100644 index 0000000..f1308ea --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/FlowableExampleActivity.java @@ -0,0 +1,84 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import org.reactivestreams.Subscriber; +import org.reactivestreams.Subscription; + +import io.reactivex.Flowable; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class FlowableExampleActivity extends AbsExampleActivity { + + private Subscription mSubscription; + + /* + * simple example using Flowable + */ + protected void doSomeWork() { + + if (mSubscription != null) { + mSubscription.cancel(); + } + + Flowable.range(0, 20) + .subscribeOn(Schedulers.io()) + .observeOn(AndroidSchedulers.mainThread()) + .subscribeWith(new Subscriber() { + + //当订阅后,会首先调用这个方法,其实就相当于onStart(), + //传入的Subscription s参数可以用于请求数据或者取消订阅 + @Override + public void onSubscribe(Subscription s) { + FlowableExampleActivity.this.onSubscribe("",null); + mSubscription = s; + //要说明一下,request这个方法若不调用,下游的onNext与OnComplete都不会调用; + // 若你写的数量小于真实数据量,只会传你的个数,而且不会调用onComplete方法(毕竟没有传完嘛,当然没有complete) +// mSubscription.request(10); + mSubscription.request(1); + } + + @Override + public void onNext(Integer o) { +// SystemClock.sleep(1000); + FlowableExampleActivity.this.onNext("",o); + mSubscription.request(1); + } + + @Override + public void onError(Throwable t) { + FlowableExampleActivity.this.onError("",t); + } + + @Override + public void onComplete() { + FlowableExampleActivity.this.onComplete(""); + } + }); + + +// Flowable observable = Flowable.just(1, 2, 3, 4); +// +// observable.reduce(50, new BiFunction() { +// @Override +// public Integer apply(Integer t1, Integer t2) { +// return t1 + t2; +// } +// }).subscribe(getObserver()); + + } + + + @Override + protected void onDestroy() { + if (mSubscription!=null) { + mSubscription.cancel(); + } + super.onDestroy(); +// mCompositeDisposable.dispose(); + } +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/IntervalExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/IntervalExampleActivity.java new file mode 100644 index 0000000..5c2bc7c --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/IntervalExampleActivity.java @@ -0,0 +1,43 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import java.util.concurrent.TimeUnit; + +import io.reactivex.Observable; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.disposables.CompositeDisposable; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class IntervalExampleActivity extends AbsExampleActivity { + + private final CompositeDisposable disposables = new CompositeDisposable(); + + @Override + protected void onDestroy() { + super.onDestroy(); + disposables.clear(); // clearing it : do not emit after destroy + } + + /* + * simple example using interval to run task at an interval of 2 sec + * which start immediately + */ + protected void doSomeWork() { + disposables.add(getObservable() + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribeWith(this.getDisposableObserver())); + } + + private Observable getObservable() { + //第一个参数是initialDelay初始化延迟,第二个参数是间隔时间。 + return Observable.interval(0, 1, TimeUnit.SECONDS); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/LastOperatorExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/LastOperatorExampleActivity.java new file mode 100644 index 0000000..abc2fe9 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/LastOperatorExampleActivity.java @@ -0,0 +1,24 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; + +/** + * Created by techteam on 13/09/16. + */ +public class LastOperatorExampleActivity extends AbsExampleActivity { + + protected void doSomeWork() { + getObservable().last("ADefault") // the default item ("ADefault") to emit if the source ObservableSource is empty + .subscribe(this.getSingleObserver()); + //经过Last操作后,Observable变成SingleObservable了 + } + + private Observable getObservable() { +// return Observable.empty(); + return Observable.just("A1", "A2", "A3", "A4", "A5", "A6"); + } + + +} diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/MapExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/MapExampleActivity.java new file mode 100644 index 0000000..9d9eb0e --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/MapExampleActivity.java @@ -0,0 +1,56 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; +import com.rxjava2.android.samples.model.ApiUser; +import com.rxjava2.android.samples.model.User; +import com.rxjava2.android.samples.utils.Utils; + +import java.util.List; + +import io.reactivex.Observable; +import io.reactivex.ObservableEmitter; +import io.reactivex.ObservableOnSubscribe; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.functions.Function; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class MapExampleActivity extends AbsExampleActivity { + + /* + * Here we are getting ApiUser Object from api server + * then we are converting it into User Object because + * may be our database support User Not ApiUser Object + * Here we are using Map Operator to do that + */ + protected void doSomeWork() { + getObservable() + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .map(new Function, List>() { + + @Override + public List apply(List apiUsers) throws Exception { + return Utils.convertApiUserListToUserList(apiUsers); + } + }) + .subscribe(getObserver()); + } + + private Observable> getObservable() { + return Observable.create(new ObservableOnSubscribe>() { + @Override + public void subscribe(ObservableEmitter> e) throws Exception { + if (!e.isDisposed()) { + e.onNext(Utils.getApiUserList()); + e.onComplete(); + } + } + }); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/MergeExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/MergeExampleActivity.java new file mode 100644 index 0000000..e75effe --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/MergeExampleActivity.java @@ -0,0 +1,28 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; + +/** + * Created by amitshekhar on 28/08/16. + */ +public class MergeExampleActivity extends AbsExampleActivity { + /* + * Using merge operator to combine Observable : merge does not maintain + * the order of Observable. + * It will emit all the 7 values may not be in order + * Ex - "A1", "B1", "A2", "A3", "A4", "B2", "B3" - may be anything + */ + protected void doSomeWork() { + final String[] aStrings = {"A1", "A2", "A3", "A4"}; + final String[] bStrings = {"B1", "B2", "B3"}; + + final Observable aObservable = Observable.fromArray(aStrings); + final Observable bObservable = Observable.fromArray(bStrings); + + Observable.merge(aObservable, bObservable) + .subscribe(getObserver()); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/PublishSubjectExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/PublishSubjectExampleActivity.java new file mode 100644 index 0000000..83c3a09 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/PublishSubjectExampleActivity.java @@ -0,0 +1,41 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.subjects.PublishSubject; + +/** + * Created by amitshekhar on 17/12/16. + */ + +public class PublishSubjectExampleActivity extends AbsExampleActivity { + + /* PublishSubject emits to an observer only those items that are emitted + * by the source Observable, subsequent to the time of the subscription. + */ + protected void doSomeWork() { + + //PublishSubject就像鼠标事件,不管有没有订阅者,他都按照他的设定发事件,什么时候有订阅者,那么订阅者就从那个时候获取事件。 + //与之相反的是ReplaySubject,不管订阅者什么时候订阅,都能获取完整事件。 + PublishSubject source = PublishSubject.create(); + source.onNext(-1); + source.onNext(0); + + source.subscribe(this.getObserver("First")); // it will get 1, 2, 3, 4 and onComplete + + source.onNext(1); + source.onNext(2); + source.onNext(3); +// source.onComplete();//如果在这里onComplete了,那么后面的订阅者只能收到onComplete事件 + + /* + * it will emit 4 and onComplete for second observer also. + */ + source.subscribe(this.getObserver("Second")); + + source.onNext(4); + source.onComplete(); + + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/ReduceExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/ReduceExampleActivity.java new file mode 100644 index 0000000..6c1fa8c --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/ReduceExampleActivity.java @@ -0,0 +1,31 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; +import io.reactivex.functions.BiFunction; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class ReduceExampleActivity extends AbsExampleActivity { + /* + * simple example using reduce to add all the number + */ + protected void doSomeWork() { + getObservable() + .reduce(new BiFunction() { + @Override + public Integer apply(Integer t1, Integer t2) { + return t1 + t2; + } + }) + .subscribe(this.getMaybeObserver()); + } + + private Observable getObservable() { + return Observable.just(1, 2, 3, 4); + } + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/ReplayExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/ReplayExampleActivity.java new file mode 100644 index 0000000..3d3e4f7 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/ReplayExampleActivity.java @@ -0,0 +1,45 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.observables.ConnectableObservable; +import io.reactivex.subjects.PublishSubject; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class ReplayExampleActivity extends AbsExampleActivity { + + /* Using replay operator, replay ensure that all observers see the same sequence + * of emitted items, even if they subscribe after the Observable has begun emitting items + */ + protected void doSomeWork() { + + PublishSubject source = PublishSubject.create(); + ConnectableObservable connectableObservable = source.replay(3); // bufferSize = 3 to retain 3 values to replay + connectableObservable.connect(); // connecting the connectableObservable + + source.onNext(-1); + source.onNext(0); + connectableObservable.subscribe(this.getObserver("First")); + + source.onNext(1); + source.onNext(2); + source.onNext(3); + source.onNext(4); + source.onComplete(); + + + /* + * Replay操作符会给onComplete后订阅Observable的订阅者Replay发射最后几个元素。 + * 在onComplete之前订阅的不受影响(会收到完整的元素。哪怕在订阅之前已经开始onNext数据了) + * + * it will emit 2, 3, 4 as (count = 3), retains the 3 values for replay + */ + connectableObservable.subscribe(this.getObserver("Second")); + + } + + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/ReplaySubjectExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/ReplaySubjectExampleActivity.java new file mode 100644 index 0000000..484b7e8 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/ReplaySubjectExampleActivity.java @@ -0,0 +1,49 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.subjects.ReplaySubject; + +/** + * Created by amitshekhar on 17/12/16. + */ + +public class ReplaySubjectExampleActivity extends AbsExampleActivity { + + /* ReplaySubject emits to any observer all of the items that were emitted + * by the source Observable, regardless of when the observer subscribes. + */ + protected void doSomeWork() { + + //ReplaySubject 和 PublishSubject相反, + //ReplaySubject不管订阅者什么时候订阅都能获取到完整的发射数据。 + //而PublishSubject会一直按照自己的步调发射数据,你在哪订阅就从这个时间点开始才能获取到事件 + //所谓的完整数据指的是从第一个onNext 一直到 onComplete 或 onError + ReplaySubject source = ReplaySubject.create(); + + + source.onNext(-1); + source.onNext(0); + + source.subscribe(this.getObserver("First")); // it will get -1, 0, 1, 2, 3, 4 + + source.onNext(1); + source.onNext(2); + source.onNext(3); + source.onNext(4); +// source.onComplete(); +//OnComplete后发射数据就无效了。如果仍要发射数据,需要重新创建ReplaySubject 并重新订阅 +// source.onNext(5); + + /* + * it will emit -1, 0, 1, 2, 3, 4 for second observer also as we have used replay + */ + source.subscribe(this.getObserver("Second")); + + source.onNext(6); + source.onComplete(); + + } + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/ScanExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/ScanExampleActivity.java new file mode 100644 index 0000000..e9ee9dd --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/ScanExampleActivity.java @@ -0,0 +1,35 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.functions.BiFunction; +import io.reactivex.schedulers.Schedulers; + +public class ScanExampleActivity extends AbsExampleActivity { + + /* Using scan operator, it sends also the previous result */ + protected void doSomeWork() { + getObservable() + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .scan(new BiFunction() { + @Override + public Integer apply(Integer int1, Integer int2) throws Exception { + return int1 + int2; + } + }) + .subscribe(getObserver()); + //Scan又叫累加器。将原始第一个与第二个应用函数的值作为第二个发射出去数据(第一个发射的数据就是原始第一个数据) + //第三个发射的数据是原始第三个与第二个发射出去的应用函数后的值。 + + + } + + private Observable getObservable() { + return Observable.just(1, 2, 3, 4, 5); + } +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/SimpleExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/SimpleExampleActivity.java new file mode 100644 index 0000000..c6168ce --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/SimpleExampleActivity.java @@ -0,0 +1,31 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class SimpleExampleActivity extends AbsExampleActivity { + + /* + * simple example to emit two value one by one + */ + protected void doSomeWork() { + getObservable() + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribe(getObserver()); + } + + private Observable getObservable() { + return Observable.just("Cricket", "Football"); + } + + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/SingleObserverExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/SingleObserverExampleActivity.java new file mode 100644 index 0000000..87b11b3 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/SingleObserverExampleActivity.java @@ -0,0 +1,51 @@ +package com.rxjava2.android.samples.operators; + +import android.util.Log; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Single; +import io.reactivex.SingleEmitter; +import io.reactivex.SingleOnSubscribe; +import io.reactivex.functions.Consumer; +import io.reactivex.functions.Function; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class SingleObserverExampleActivity extends AbsExampleActivity { + + /* + * simple example using SingleObserver + */ + protected void doSomeWork() { + Single.create(new SingleOnSubscribe() { + @Override + public void subscribe(SingleEmitter e) throws Exception { + if (!e.isDisposed()) { + e.onSuccess("Hello Success!"); +// e.onError(new RuntimeException("Occur RuntimeException")); +// throw new RuntimeException("Occur RuntimeException"); + } + } + }).doOnSuccess(new Consumer() { + @Override + public void accept(String s) throws Exception { + Log.d(TAG, "doOnSuccess: " + s); + } + }).doOnError(new Consumer() { + @Override + public void accept(Throwable throwable) throws Exception { + Log.e(TAG, "doOnError: " + throwable.getMessage()); + } + }).onErrorReturn(new Function() { + @Override + public String apply(Throwable throwable) throws Exception { + return "Exception message: "+throwable.getMessage(); + } + }).subscribe(getSingleObserver()); +// Single.just("Amit") +// .subscribe(getSingleObserver()); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/SkipExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/SkipExampleActivity.java new file mode 100644 index 0000000..aecedc2 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/SkipExampleActivity.java @@ -0,0 +1,31 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class SkipExampleActivity extends AbsExampleActivity { + + /* Using skip operator, it only not emit + * the first 2 values. + */ + protected void doSomeWork() { + getObservable() + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .skip(2) + .subscribe(getObserver()); + } + + private Observable getObservable() { + return Observable.just(1, 2, 3, 4, 5); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/TakeExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/TakeExampleActivity.java new file mode 100644 index 0000000..168de2d --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/TakeExampleActivity.java @@ -0,0 +1,32 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import io.reactivex.Observable; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class TakeExampleActivity extends AbsExampleActivity { + + /* + * Using take operator, it only emits + * required number of values. here only 3 out of 5 + */ + protected void doSomeWork() { + getObservable() + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .take(3) + .subscribe(getObserver()); + } + + private Observable getObservable() { + return Observable.just(1, 2, 3, 4, 5); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/ThrottleFirstExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/ThrottleFirstExampleActivity.java new file mode 100644 index 0000000..bc05cb1 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/ThrottleFirstExampleActivity.java @@ -0,0 +1,58 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import java.util.concurrent.TimeUnit; + +import io.reactivex.Observable; +import io.reactivex.ObservableEmitter; +import io.reactivex.ObservableOnSubscribe; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + + +/** + * Created by threshold on 2017/1/11. + */ + +public class ThrottleFirstExampleActivity extends AbsExampleActivity { + + protected void doSomeWork() { + getObservable() + .throttleFirst(500, TimeUnit.MILLISECONDS) + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribe(getObserver()); + } + + private Observable getObservable() { + return Observable.create(new ObservableOnSubscribe() { + @Override + public void subscribe(ObservableEmitter emitter) throws Exception { + // send events with simulated time wait + Thread.sleep(0); + emitter.onNext(1); // skip + emitter.onNext(2); // deliver + Thread.sleep(505); + emitter.onNext(3); // skip + Thread.sleep(99); + emitter.onNext(4); // skip + Thread.sleep(100); + emitter.onNext(5); // skip + emitter.onNext(6); // deliver + Thread.sleep(305); + emitter.onNext(7); // deliver + Thread.sleep(510); + emitter.onComplete(); + /* + 0----------500-------------1000----------1500 + 0------------505--604--704--1009----------1519 + 1,2-----------3----4---5,6----7-----------Complete + */ + } + }); + } + +} diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/ThrottleLastExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/ThrottleLastExampleActivity.java new file mode 100644 index 0000000..33f8ba9 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/ThrottleLastExampleActivity.java @@ -0,0 +1,62 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import java.util.concurrent.TimeUnit; + +import io.reactivex.Observable; +import io.reactivex.ObservableEmitter; +import io.reactivex.ObservableOnSubscribe; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 22/12/16. + */ + +public class ThrottleLastExampleActivity extends AbsExampleActivity { + + /* + * Using throttleLast() -> emit the most recent items emitted by an Observable within + * periodic time intervals, so here it will emit 2, 6 and 7 as we have simulated it to be the + * last the element in the interval of 500 millis + */ + protected void doSomeWork() { + getObservable() + .throttleLast(500, TimeUnit.MILLISECONDS) + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribe(getObserver()); + } + + private Observable getObservable() { + return Observable.create(new ObservableOnSubscribe() { + @Override + public void subscribe(ObservableEmitter emitter) throws Exception { + // send events with simulated time wait + Thread.sleep(0); + emitter.onNext(1); // skip + emitter.onNext(2); // deliver + Thread.sleep(505); + emitter.onNext(3); // skip + Thread.sleep(99); + emitter.onNext(4); // skip + Thread.sleep(100); + emitter.onNext(5); // skip + emitter.onNext(6); // deliver + Thread.sleep(305); + emitter.onNext(7); // deliver + Thread.sleep(510); + emitter.onComplete(); + /* + 0----------500-------------1000----------1500 + 0------------505--604--704--1009----------1519 + 1,2-----------3----4---5,6----7-----------Complete + */ + } + }); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/TimerExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/TimerExampleActivity.java new file mode 100644 index 0000000..f4244ee --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/TimerExampleActivity.java @@ -0,0 +1,33 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; + +import java.util.concurrent.TimeUnit; + +import io.reactivex.Observable; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class TimerExampleActivity extends AbsExampleActivity { + + /* + * simple example using timer to do something after 2 second + */ + protected void doSomeWork() { + getObservable() + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribe(getObserver()); + } + + private Observable getObservable() { + //延迟2秒后发射一个0. + // timer就是定时器,到点后仅发射一个0 + return Observable.timer(2, TimeUnit.SECONDS); + } +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/WindowExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/WindowExampleActivity.java new file mode 100644 index 0000000..37592f2 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/WindowExampleActivity.java @@ -0,0 +1 @@ +package com.rxjava2.android.samples.operators; import android.util.Log; import com.rxjava2.android.samples.AbsExampleActivity; import com.rxjava2.android.samples.utils.AppConstant; import java.util.concurrent.TimeUnit; import io.reactivex.Observable; import io.reactivex.android.schedulers.AndroidSchedulers; import io.reactivex.functions.Consumer; import io.reactivex.schedulers.Schedulers; public class WindowExampleActivity extends AbsExampleActivity { /* * Sample of using window */ protected void doSomeWork() { // Observable.create(new ObservableOnSubscribe() { // @Override // public void subscribe(ObservableEmitter e) throws Exception { // //e.onNext(-1L); //在这里发射一个就相当于0秒也发射了一个数据 // for (int i =0 ;i<12;i++) { // Thread.sleep(1000); // e.onNext((long)i); // Log.d(TAG, "e.onNext:" + i); // } // e.onComplete(); // } // }) Observable.interval(1,TimeUnit.SECONDS).take(12) .window(3, TimeUnit.SECONDS) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(new Consumer>() { @Override public void accept(Observable observable) { Log.d(TAG, "Sub Divide begin...."); textView.append("Sub Divide begin ...."); textView.append(AppConstant.LINE_SEPARATOR); observable .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(new Consumer() { @Override public void accept(Long aLong) { Log.d(TAG, "Next:" + aLong); textView.append("Next:" + aLong); textView.append(AppConstant.LINE_SEPARATOR); } }); } }); /* 输出 subdivide begin…… Next:0 Next:1 subdivide begin…… Next:2 Next:3 Next:4 subdivide begin…… Next:5 Next:6 Next:7 subdivide begin…… Next:8 Next:9 Next:10 subdivide begin…… Next:11 为什么不是 0,1,2 3,4,5 ..... 这种形式呢。 同学,问的非常好!!! 因为0秒的时刻啥也没发射,1秒的时刻发射了0,2秒的时刻发射了1,3秒的时刻发射了2 "对呀,前3秒嘛,因该0,1,2嘛",但是3秒是指时间片段,[0,1) [1,2) [2,3) 左闭右开 这三秒只发射了0和1 */ } } \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/operators/ZipExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/operators/ZipExampleActivity.java new file mode 100644 index 0000000..58e504a --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/operators/ZipExampleActivity.java @@ -0,0 +1,65 @@ +package com.rxjava2.android.samples.operators; + +import com.rxjava2.android.samples.AbsExampleActivity; +import com.rxjava2.android.samples.model.User; +import com.rxjava2.android.samples.utils.Utils; + +import java.util.List; + +import io.reactivex.Observable; +import io.reactivex.ObservableEmitter; +import io.reactivex.ObservableOnSubscribe; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.functions.BiFunction; +import io.reactivex.schedulers.Schedulers; + +/** + * Created by amitshekhar on 27/08/16. + */ +public class ZipExampleActivity extends AbsExampleActivity { + /* + * Here we are getting two user list + * One, the list of cricket fans + * Another one, the list of football fans + * Then we are finding the list of users who loves both + */ + protected void doSomeWork() { + Observable.zip(getCricketFansObservable(), getFootballFansObservable(), + new BiFunction, List, List>() { + @Override + public List apply(List cricketFans, List footballFans) throws Exception { + return Utils.filterUserWhoLovesBoth(cricketFans, footballFans); + } + }) + // Run on a background thread + .subscribeOn(Schedulers.io()) + // Be notified on the main thread + .observeOn(AndroidSchedulers.mainThread()) + .subscribe(getObserver()); + } + + private Observable> getCricketFansObservable() { + return Observable.create(new ObservableOnSubscribe>() { + @Override + public void subscribe(ObservableEmitter> e) throws Exception { + if (!e.isDisposed()) { + e.onNext(Utils.getUserListWhoLovesCricket()); + e.onComplete(); + } + } + }); + } + + private Observable> getFootballFansObservable() { + return Observable.create(new ObservableOnSubscribe>() { + @Override + public void subscribe(ObservableEmitter> e) throws Exception { + if (!e.isDisposed()) { + e.onNext(Utils.getUserListWhoLovesFootball()); + e.onComplete(); + } + } + }); + } + +} \ No newline at end of file diff --git a/app/src/main/java/com/rxjava2/android/samples/utils/Utils.java b/app/src/main/java/com/rxjava2/android/samples/utils/Utils.java index 0311e64..ddf699f 100644 --- a/app/src/main/java/com/rxjava2/android/samples/utils/Utils.java +++ b/app/src/main/java/com/rxjava2/android/samples/utils/Utils.java @@ -1,10 +1,22 @@ package com.rxjava2.android.samples.utils; +import android.util.Log; + import com.rxjava2.android.samples.model.ApiUser; import com.rxjava2.android.samples.model.User; import java.util.ArrayList; import java.util.List; +import java.util.Random; +import java.util.function.Consumer; + +import io.reactivex.Observable; +import io.reactivex.Observer; +import io.reactivex.disposables.Disposable; +import io.reactivex.functions.BiFunction; + +import static android.icu.lang.UCharacter.GraphemeClusterBreak.T; +import static com.rxjava2.android.samples.R.id.textView; /** * Created by amitshekhar on 27/08/16. @@ -120,4 +132,5 @@ public static List filterUserWhoLovesBoth(List cricketFans, List +