diff --git a/README.md b/README.md index 1e1d05a..8aa4e01 100644 --- a/README.md +++ b/README.md @@ -2,145 +2,111 @@ # Learning RxJava 2 for Android by example -[![Mindorks](https://img.shields.io/badge/mindorks-opensource-blue.svg)](https://mindorks.com/open-source-projects) -[![Mindorks Community](https://img.shields.io/badge/join-community-blue.svg)](https://mindorks.com/join-community) -[![Mindorks Android Store](https://img.shields.io/badge/Mindorks%20Android%20Store-RxJava2%20Android%20Samples-blue.svg?style=flat)](https://mindorks.com/android/store) -[![Open Source Love](https://badges.frapsoft.com/os/v1/open-source.svg?v=102)](https://opensource.org/licenses/Apache-2.0) -[![License](https://img.shields.io/badge/license-Apache%202.0-blue.svg)](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/LICENSE) - -### Get the complete RxJava Course [here](https://mindorks.com/course/learn-rxjava) - -## How to use RxJava 2 in Android Application -## How to migrate from RxJava 1.0 to RxJava 2.0 - ### This project is for : * who is migrating to RxJava 2 * or just started with RxJava. -### Just Build the project and start learning RxJava by examples. +## About me -RxJava 2.0 has been completely rewritten from scratch on top of the Reactive-Streams specification. The specification itself has evolved out of RxJava 1.x and provides a common baseline for reactive systems and libraries. +Hi, I am Amit Shekhar, Founder @ [Outcome School](https://outcomeschool.com) • IIT 2010-14 • I have taught and mentored many developers, and their efforts landed them high-paying tech jobs, helped many tech companies in solving their unique problems, and created many open-source libraries being used by top companies. I am passionate about sharing knowledge through open-source, blogs, and videos. + +### Follow Amit Shekhar + +- [X/Twitter](https://twitter.com/amitiitbhu) +- [LinkedIn](https://www.linkedin.com/in/amit-shekhar-iitbhu) +- [GitHub](https://github.com/amitshekhariitbhu) -Because Reactive-Streams has a different architecture, it mandates changes to some well known RxJava types. +### Follow Outcome School +- [YouTube](https://youtube.com/@OutcomeSchool) +- [X/Twitter](https://x.com/outcome_school) +- [LinkedIn](https://www.linkedin.com/company/outcomeschool) +- [GitHub](http://github.com/OutcomeSchool) -# Migration From RxJava 1.0 to RxJava 2.0 +## I teach at Outcome School -To allow having RxJava 1 and RxJava 2 side-by-side, RxJava 2 is under the maven coordinates -io.reactivex.rxjava2:rxjava:2.x.y and classes are accessible below io.reactivex. +- AI and Machine Learning +- Android -Users switching from 1.x to 2.x have to re-organize their imports, but carefully. +Join Outcome School and get a high-paying tech job: [Outcome School](https://outcomeschool.com) + +### Just Build the project and start learning RxJava by examples. + +RxJava 2.0 has been completely rewritten from scratch on top of the Reactive-Streams specification. The specification itself has evolved out of RxJava 1.x and provides a common baseline for reactive systems and libraries. ### Using RxJava 2.0 Library in your application Add this in your build.gradle ```groovy -compile 'io.reactivex.rxjava2:rxjava:2.2.2' +compile 'io.reactivex.rxjava2:rxjava:X.X.X' ``` If you are using RxAndroid also, then add the following ```groovy -compile 'io.reactivex.rxjava2:rxandroid:2.1.0' +compile 'io.reactivex.rxjava2:rxandroid:X.X.X' ``` -# RxJava 2 Examples present in this sample project - -* RxJava 2.0 Example using `CompositeDisposable` as `CompositeSubscription` and `Subscription` have -been removed. - -* RxJava 2 Example using `Flowable`. - -* RxJava 2 Example using `SingleObserver`, `CompletableObserver`. +# RxJava 2 Operators Examples present in this sample project: -* RxJava 2 Example using RxJava2 operators such as `map, zip, take, reduce, flatMap, filter, buffer, skip, merge, concat, replay`, and much more: - -* RxJava 2 Android Samples using `Function` as `Func1` has been removed. - -* RxJava 2 Android Samples using `BiFunction` as `Func2` has been removed. - -* And many others android examples - -# Quick Look on few changes done in RxJava2 over RxJava1 - -RxJava1 -> RxJava2 - -* `onCompleted` -> `onComplete` - without the trailing d -* `Func1` -> `Function` -* `Func2` -> `BiFunction` -* `CompositeSubscription` -> `CompositeDisposable` -* `limit` operator has been removed - Use `take` in RxJava2 -* and much more. - -# Operators : -* `Map` -> transform the items emitted by an Observable by applying a function to each item +* `Map` -> transform the items emitted by an Observable by applying a function to each item. Blog: [RxJava Operator Map vs FlatMap](https://outcomeschool.com/blog/rxjava-map-vs-flatmap) * `Zip` -> combine the emissions of multiple Observables together via a specified function and emit single items for each combination based on the results of this function * `Filter` -> emit only those items from an Observable that pass a predicate test -* `FlatMap` -> transform the items emitted by an Observable into Observables, then flatten the emissions from those into a single Observable -* `Take` -> emit only the first n items emitted by an Observable +* `FlatMap` -> transform the items emitted by an Observable into Observables, then flatten the emissions from those into a single Observable. Blog: [RxJava Operator Map vs FlatMap](https://outcomeschool.com/blog/rxjava-map-vs-flatmap) +* `Take` -> emit only the first n items emitted by an Observable. [Blog for reference](https://outcomeschool.com/blog/rxjava-interval-operator) * `Reduce` -> apply a function to each item emitted by an Observable, sequentially, and emit the final value * `Skip` -> suppress the first n items emitted by an Observable * `Buffer` -> periodically gather items emitted by an Observable into bundles and emit these bundles rather than emitting the items one at a time -* `Concat` -> emit the emissions from two or more Observables without interleaving them +* `Concat` -> emit the emissions from two or more Observables without interleaving them. [Blog for reference](https://outcomeschool.com/blog/rxjava-concat-operator) * `Replay` -> ensure that all observers see the same sequence of emitted items, even if they subscribe after the Observable has begun emitting items * `Merge` -> combine multiple Observables into one by merging their emissions -* `SwitchMap` -> ransform the items emitted by an Observable into Observables, and mirror those items emitted by the most-recently transformed Observable +* `SwitchMap` -> transform the items emitted by an Observable into Observables, and mirror those items emitted by the most-recently transformed Observable # Highlights of the examples : -* [DisposableExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/DisposableExampleActivity.java) - Using `CompositeDisposable` +* [DisposableExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/DisposableExampleActivity.java) - Using `CompositeDisposable`. [Blog for reference](https://outcomeschool.com/blog/dispose-vs-clear-compositedisposable-rxjava) * [FlowableExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/FlowableExampleActivity.java) - Using `Flowable` and `reduce` operator * [SingleObserverExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/SingleObserverExampleActivity.java) - Using `SingleObserver` * [CompletableObserverActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/CompletableObserverExampleActivity.java) - Using `CompletableObserver` -* [MapExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/MapExampleActivity.java) - Using `map` Operator +* [MapExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/MapExampleActivity.java) - Using `map` Operator. Blog: [RxJava Operator Map vs FlatMap](https://outcomeschool.com/blog/rxjava-map-vs-flatmap) * [ZipExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/ZipExampleActivity.java) - Using `zip` Operator * [BufferExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/BufferExampleActivity.java) - Using `buffer` Operator -* [TakeExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/TakeExampleActivity.java) - Using `take` Operator +* [TakeExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/TakeExampleActivity.java) - Using `take` Operator. [Blog for reference](https://outcomeschool.com/blog/rxjava-interval-operator) * [ReduceExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/ReduceExampleActivity.java) - Using `reduce` Operator * [FilterExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/FilterExampleActivity.java) - Using `filter` Operator * [SkipExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/SkipExampleActivity.java) - Using `skip` Operator * [ReplayExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/ReplayExampleActivity.java) - Using `replay` Operator -* [ConcatExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/ConcatExampleActivity.java) - Using `concat` Operator +* [ConcatExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/ConcatExampleActivity.java) - Using `concat` Operator. [Blog for reference](https://outcomeschool.com/blog/rxjava-concat-operator) * [MergeExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/MergeExampleActivity.java) - Using `merge` Operator * [DeferExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/DeferExampleActivity.java) - Using `defer` Observable * [SwitchMapExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/SwitchMapExampleActivity.java) - Using `switchMap` Observable -* [IntervalExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/IntervalExampleActivity.java) - Using `Interval` -* [RxBusActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/rxbus/RxBusActivity.java) - RxBus, RxJava2Bus, EventBus, RxEventBus, [Blog for reference](https://blog.mindorks.com/implementing-eventbus-with-rxjava-rxbus-e6c940a94bd8) -* [PaginationActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/pagination/PaginationActivity.java) - Pagination for loadMore in RecyclerView +* [IntervalExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/IntervalExampleActivity.java) - Using `Interval`. [Blog for reference](https://outcomeschool.com/blog/rxjava-interval-operator) +* [RxBusActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/rxbus/RxBusActivity.java) - RxBus, RxJava2Bus, EventBus, RxEventBus +* [PaginationActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/pagination/PaginationActivity.java) - Pagination for loadMore in RecyclerView. Blog: [Pagination In RecyclerView Using RxJava Operators](https://outcomeschool.com/blog/pagination-in-recyclerview-using-rxjava-operators) * [ComposeOperatorExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/compose/ComposeOperatorExampleActivity.java) - Compose operator for reusable -* [Search Implementation](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/search/SearchActivity.java) - Using `debounce`, `switchMap`, `distinctUntilChanged`, [Blog for reference](https://blog.mindorks.com/implement-search-using-rxjava-operators-c8882b64fe1d) -* [PublishSubjectExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/PublishSubjectExampleActivity.java) - -### TODO - -* Many examples are to be added +* [Search Implementation](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/search/SearchActivity.java) - Using `debounce`, `switchMap`, `distinctUntilChanged` +* [Implement Caching Using RxJava Operators](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/cache/CacheExampleActivity.java) - Using `concat`, `firstElement` +* [PublishSubjectExampleActivity](https://github.com/amitshekhariitbhu/RxJava2-Android-Samples/blob/master/app/src/main/java/com/rxjava2/android/samples/ui/operators/PublishSubjectExampleActivity.java). Blog: [RxJava Subject - Publish, Replay, Behavior, and Async](https://outcomeschool.com/blog/rxjava-subject-publish-replay-behavior-async) ### Find this project useful ? :heart: * Support it by clicking the :star: button on the upper right of this page. :v: -### Check out an awesome MVP architecture based project which uses RxJava2, Dagger2. -* [Android-MVP-Architecture](https://github.com/MindorksOpenSource/android-mvp-architecture) - -### Check out an awesome Kotlin MVP architecture based project which uses RxJava2, Dagger2. -* [Android-Kotlin-MVP-Architecture](https://github.com/MindorksOpenSource/android-kotlin-mvp-architecture) +Thanks -### Check out an awesome library for fast and simple networking in Android. -* [Fast Android Networking Library](https://github.com/amitshekhariitbhu/Fast-Android-Networking) +**Amit Shekhar**\ +Co-Founder @ [Outcome School](https://outcomeschool.com) -### Another awesome library for debugging databases and shared preferences. -* [Android Debug Database](https://github.com/amitshekhariitbhu/Android-Debug-Database) +You can connect with me on: -### [Check out Mindorks awesome open source projects here](https://mindorks.com/open-source-projects) - -### Contact - Let's become friend - [Twitter](https://twitter.com/amitiitbhu) -- [Github](https://github.com/amitshekhariitbhu) -- [Medium](https://medium.com/@amitshekhar) +- [LinkedIn](https://www.linkedin.com/in/amit-shekhar-iitbhu) +- [GitHub](https://github.com/amitshekhariitbhu) - [Facebook](https://www.facebook.com/amit.shekhar.iitbhu) +[**Read all of our blogs here.**](https://outcomeschool.com/blog) + ### License ``` - Copyright (C) 2016 Amit Shekhar - Copyright (C) 2011 Android Open Source Project + Copyright (C) 2024 Amit Shekhar Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. @@ -157,5 +123,4 @@ RxJava1 -> RxJava2 ### Contributing to RxJava 2 Android Samples Just make pull request. You are in! - diff --git a/app/src/main/AndroidManifest.xml b/app/src/main/AndroidManifest.xml index 492a93f..720343a 100644 --- a/app/src/main/AndroidManifest.xml +++ b/app/src/main/AndroidManifest.xml @@ -24,6 +24,9 @@ + @@ -126,7 +129,8 @@ - memory = dataSource.getDataFromMemory(); + Observable disk = dataSource.getDataFromDisk(); + Observable network = dataSource.getDataFromNetwork(); + + Observable.concat(memory, disk, network) + .firstElement() + .subscribeOn(Schedulers.io()) + .observeOn(AndroidSchedulers.mainThread()) + .toObservable() + .subscribe(getObserver()); + } + + private Observer getObserver() { + return new Observer() { + + @Override + public void onSubscribe(Disposable d) { + Log.d(TAG, " onSubscribe : " + d.isDisposed()); + } + + @Override + public void onNext(Data data) { + textView.append(" onNext : " + data.source); + textView.append(AppConstant.LINE_SEPARATOR); + Log.d(TAG, " onNext : " + data.source); + } + + @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/ui/cache/model/Data.java b/app/src/main/java/com/rxjava2/android/samples/ui/cache/model/Data.java new file mode 100644 index 0000000..477144e --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/ui/cache/model/Data.java @@ -0,0 +1,12 @@ +package com.rxjava2.android.samples.ui.cache.model; + +public class Data { + + public String source; + + @SuppressWarnings("CloneDoesntDeclareCloneNotSupportedException") + @Override + public Data clone() { + return new Data(); + } +} diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/DataSource.java b/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/DataSource.java new file mode 100644 index 0000000..23403ff --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/DataSource.java @@ -0,0 +1,41 @@ +package com.rxjava2.android.samples.ui.cache.source; + +import com.rxjava2.android.samples.ui.cache.model.Data; + +import io.reactivex.Observable; + +/** + * The DataSource to handle 3 data sources - memory, disk, network + */ +public class DataSource { + + private final MemoryDataSource memoryDataSource; + private final DiskDataSource diskDataSource; + private final NetworkDataSource networkDataSource; + + public DataSource(MemoryDataSource memoryDataSource, + DiskDataSource diskDataSource, + NetworkDataSource networkDataSource) { + this.memoryDataSource = memoryDataSource; + this.diskDataSource = diskDataSource; + this.networkDataSource = networkDataSource; + } + + public Observable getDataFromMemory() { + return memoryDataSource.getData(); + } + + public Observable getDataFromDisk() { + return diskDataSource.getData().doOnNext(data -> + memoryDataSource.cacheInMemory(data) + ); + } + + public Observable getDataFromNetwork() { + return networkDataSource.getData().doOnNext(data -> { + diskDataSource.saveToDisk(data); + memoryDataSource.cacheInMemory(data); + }); + } + +} diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/DiskDataSource.java b/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/DiskDataSource.java new file mode 100644 index 0000000..29f5bfa --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/DiskDataSource.java @@ -0,0 +1,28 @@ +package com.rxjava2.android.samples.ui.cache.source; + +import com.rxjava2.android.samples.ui.cache.model.Data; + +import io.reactivex.Observable; + +/** + * Class to simulate Disk DataSource + */ +public class DiskDataSource { + + private Data data; + + public Observable getData() { + return Observable.create(emitter -> { + if (data != null) { + emitter.onNext(data); + } + emitter.onComplete(); + }); + } + + public void saveToDisk(Data data) { + this.data = data.clone(); + this.data.source = "disk"; + } + +} diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/MemoryDataSource.java b/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/MemoryDataSource.java new file mode 100644 index 0000000..2742f26 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/MemoryDataSource.java @@ -0,0 +1,28 @@ +package com.rxjava2.android.samples.ui.cache.source; + +import com.rxjava2.android.samples.ui.cache.model.Data; + +import io.reactivex.Observable; + +/** + * Class to simulate InMemory DataSource + */ +public class MemoryDataSource { + + private Data data; + + public Observable getData() { + return Observable.create(emitter -> { + if (data != null) { + emitter.onNext(data); + } + emitter.onComplete(); + }); + } + + public void cacheInMemory(Data data) { + this.data = data.clone(); + this.data.source = "memory"; + } + +} diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/NetworkDataSource.java b/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/NetworkDataSource.java new file mode 100644 index 0000000..d903fb0 --- /dev/null +++ b/app/src/main/java/com/rxjava2/android/samples/ui/cache/source/NetworkDataSource.java @@ -0,0 +1,22 @@ +package com.rxjava2.android.samples.ui.cache.source; + +import com.rxjava2.android.samples.ui.cache.model.Data; + +import io.reactivex.Observable; + + +/** + * Class to simulate Network DataSource + */ +public class NetworkDataSource { + + public Observable getData() { + return Observable.create(emitter -> { + Data data = new Data(); + data.source = "network"; + emitter.onNext(data); + emitter.onComplete(); + }); + } + +} diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/networking/NetworkingActivity.java b/app/src/main/java/com/rxjava2/android/samples/ui/networking/NetworkingActivity.java index a425a96..6878e3a 100644 --- a/app/src/main/java/com/rxjava2/android/samples/ui/networking/NetworkingActivity.java +++ b/app/src/main/java/com/rxjava2/android/samples/ui/networking/NetworkingActivity.java @@ -5,6 +5,8 @@ import android.util.Pair; import android.view.View; +import androidx.appcompat.app.AppCompatActivity; + import com.rx2androidnetworking.Rx2AndroidNetworking; import com.rxjava2.android.samples.R; import com.rxjava2.android.samples.model.ApiUser; @@ -15,7 +17,6 @@ import java.util.ArrayList; import java.util.List; -import androidx.appcompat.app.AppCompatActivity; import io.reactivex.Observable; import io.reactivex.ObservableSource; import io.reactivex.Observer; @@ -93,23 +94,25 @@ public void onComplete() { private Observable> getCricketFansObservable() { return Rx2AndroidNetworking.get("https://fierce-cove-29863.herokuapp.com/getAllCricketFans") .build() - .getObjectListObservable(User.class); + .getObjectListObservable(User.class) + .subscribeOn(Schedulers.io()); } /* - * This observable return the list of User who loves Football - */ + * This observable return the list of User who loves Football + */ private Observable> getFootballFansObservable() { return Rx2AndroidNetworking.get("https://fierce-cove-29863.herokuapp.com/getAllFootballFans") .build() - .getObjectListObservable(User.class); + .getObjectListObservable(User.class) + .subscribeOn(Schedulers.io()); } /* - * This do the complete magic, make both network call - * and then returns the list of user who loves both - * Using zip operator to get both response at a time - */ + * This do the complete magic, make both network call + * and then returns the list of user who loves both + * Using zip operator to get both response at a time + */ private void findUsersWhoLovesBoth() { // here we are using zip operator to combine both request Observable.zip(getCricketFansObservable(), getFootballFansObservable(), diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/operators/WindowExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/ui/operators/WindowExampleActivity.java index db2ed70..881db12 100644 --- a/app/src/main/java/com/rxjava2/android/samples/ui/operators/WindowExampleActivity.java +++ b/app/src/main/java/com/rxjava2/android/samples/ui/operators/WindowExampleActivity.java @@ -1 +1,77 @@ -package com.rxjava2.android.samples.ui.operators; import android.os.Bundle; import android.util.Log; import android.view.View; import android.widget.Button; import android.widget.TextView; import com.rxjava2.android.samples.R; import com.rxjava2.android.samples.utils.AppConstant; import java.util.concurrent.TimeUnit; import androidx.appcompat.app.AppCompatActivity; import io.reactivex.Observable; import io.reactivex.android.schedulers.AndroidSchedulers; import io.reactivex.functions.Consumer; import io.reactivex.schedulers.Schedulers; public class WindowExampleActivity extends AppCompatActivity { private static final String TAG = WindowExampleActivity.class.getSimpleName(); Button btn; TextView textView; @Override protected void onCreate(Bundle savedInstanceState) { super.onCreate(savedInstanceState); setContentView(R.layout.activity_example); btn = findViewById(R.id.btn); textView = findViewById(R.id.textView); btn.setOnClickListener(new View.OnClickListener() { @Override public void onClick(View view) { doSomeWork(); } }); } /* * Example using window operator -> It periodically * subdivide items from an Observable into * Observable windows and emit these windows rather than * emitting the items one at a time */ protected void doSomeWork() { Observable.interval(1, TimeUnit.SECONDS).take(12) .window(3, TimeUnit.SECONDS) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(getConsumer()); } public Consumer> getConsumer() { return 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 value) { Log.d(TAG, "Next:" + value); textView.append("Next:" + value); textView.append(AppConstant.LINE_SEPARATOR); } }); } }; } } \ No newline at end of file +package com.rxjava2.android.samples.ui.operators; + +import android.os.Bundle; +import android.util.Log; +import android.view.View; +import android.widget.Button; +import android.widget.TextView; + +import com.rxjava2.android.samples.R; +import com.rxjava2.android.samples.utils.AppConstant; + +import java.util.concurrent.TimeUnit; + +import androidx.appcompat.app.AppCompatActivity; +import io.reactivex.Observable; +import io.reactivex.android.schedulers.AndroidSchedulers; +import io.reactivex.functions.Consumer; +import io.reactivex.schedulers.Schedulers; + +public class WindowExampleActivity extends AppCompatActivity { + + private static final String TAG = WindowExampleActivity.class.getSimpleName(); + Button btn; + TextView textView; + + @Override + protected void onCreate(Bundle savedInstanceState) { + super.onCreate(savedInstanceState); + setContentView(R.layout.activity_example); + btn = findViewById(R.id.btn); + textView = findViewById(R.id.textView); + + btn.setOnClickListener(new View.OnClickListener() { + @Override + public void onClick(View view) { + doSomeWork(); + } + }); + } + + /* + * Example using window operator -> It periodically + * subdivide items from an Observable into + * Observable windows and emit these windows rather than + * emitting the items one at a time + */ + protected void doSomeWork() { + + Observable.interval(1, TimeUnit.SECONDS).take(12) + .window(3, TimeUnit.SECONDS) + .subscribeOn(Schedulers.io()) + .observeOn(AndroidSchedulers.mainThread()) + .subscribe(getConsumer()); + } + + public Consumer> getConsumer() { + return 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 value) { + Log.d(TAG, "Next:" + value); + textView.append("Next:" + value); + textView.append(AppConstant.LINE_SEPARATOR); + } + }); + } + }; + } +} diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/operators/ZipExampleActivity.java b/app/src/main/java/com/rxjava2/android/samples/ui/operators/ZipExampleActivity.java index 18c97f3..11cb394 100644 --- a/app/src/main/java/com/rxjava2/android/samples/ui/operators/ZipExampleActivity.java +++ b/app/src/main/java/com/rxjava2/android/samples/ui/operators/ZipExampleActivity.java @@ -48,11 +48,11 @@ public void onClick(View view) { } /* - * 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 - */ + * 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>() { @@ -77,7 +77,7 @@ public void subscribe(ObservableEmitter> e) { e.onComplete(); } } - }); + }).subscribeOn(Schedulers.io()); } private Observable> getFootballFansObservable() { @@ -89,7 +89,7 @@ public void subscribe(ObservableEmitter> e) { e.onComplete(); } } - }); + }).subscribeOn(Schedulers.io()); } private Observer> getObserver() { diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/pagination/PaginationActivity.java b/app/src/main/java/com/rxjava2/android/samples/ui/pagination/PaginationActivity.java index 2cca2d7..e9c89a9 100644 --- a/app/src/main/java/com/rxjava2/android/samples/ui/pagination/PaginationActivity.java +++ b/app/src/main/java/com/rxjava2/android/samples/ui/pagination/PaginationActivity.java @@ -4,25 +4,22 @@ import android.view.View; import android.widget.ProgressBar; -import com.rxjava2.android.samples.R; +import androidx.appcompat.app.AppCompatActivity; +import androidx.recyclerview.widget.LinearLayoutManager; +import androidx.recyclerview.widget.RecyclerView; -import org.reactivestreams.Publisher; +import com.rxjava2.android.samples.R; import java.util.ArrayList; import java.util.List; import java.util.concurrent.TimeUnit; -import androidx.appcompat.app.AppCompatActivity; -import androidx.recyclerview.widget.LinearLayoutManager; -import androidx.recyclerview.widget.RecyclerView; -import io.reactivex.Flowable; +import io.reactivex.Single; import io.reactivex.android.schedulers.AndroidSchedulers; -import io.reactivex.annotations.NonNull; import io.reactivex.disposables.CompositeDisposable; import io.reactivex.disposables.Disposable; -import io.reactivex.functions.Consumer; -import io.reactivex.functions.Function; import io.reactivex.processors.PublishProcessor; +import io.reactivex.schedulers.Schedulers; /** * Created by amitshekhar on 15/03/17. @@ -94,23 +91,23 @@ private void subscribeForData() { Disposable disposable = paginator .onBackpressureDrop() - .concatMap(new Function>>() { - @Override - public Publisher> apply(@NonNull Integer page) { - loading = true; - progressBar.setVisibility(View.VISIBLE); - return dataFromNetwork(page); - } + .doOnNext(page -> { + loading = true; + progressBar.setVisibility(View.VISIBLE); }) + .concatMapSingle(page -> dataFromNetwork(page) + .subscribeOn(Schedulers.io()) + .doOnError(throwable -> { + // handle error + }) + // continue emission in case of error also + .onErrorReturn(throwable -> new ArrayList<>())) .observeOn(AndroidSchedulers.mainThread()) - .subscribe(new Consumer>() { - @Override - public void accept(@NonNull List items) { - paginationAdapter.addItems(items); - paginationAdapter.notifyDataSetChanged(); - loading = false; - progressBar.setVisibility(View.INVISIBLE); - } + .subscribe(items -> { + paginationAdapter.addItems(items); + paginationAdapter.notifyDataSetChanged(); + loading = false; + progressBar.setVisibility(View.INVISIBLE); }); compositeDisposable.add(disposable); @@ -122,18 +119,15 @@ public void accept(@NonNull List items) { /** * Simulation of network data */ - private Flowable> dataFromNetwork(final int page) { - return Flowable.just(true) + private Single> dataFromNetwork(final int page) { + return Single.just(true) .delay(2, TimeUnit.SECONDS) - .map(new Function>() { - @Override - public List apply(@NonNull Boolean value) { - List items = new ArrayList<>(); - for (int i = 1; i <= 10; i++) { - items.add("Item " + (page * 10 + i)); - } - return items; + .map(value -> { + List items = new ArrayList<>(); + for (int i = 1; i <= 10; i++) { + items.add("Item " + (page * 10 + i)); } + return items; }); } } diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/search/RxSearchObservable.java b/app/src/main/java/com/rxjava2/android/samples/ui/search/RxSearchObservable.java index 6360030..9189214 100644 --- a/app/src/main/java/com/rxjava2/android/samples/ui/search/RxSearchObservable.java +++ b/app/src/main/java/com/rxjava2/android/samples/ui/search/RxSearchObservable.java @@ -22,7 +22,7 @@ public static Observable fromView(SearchView searchView) { searchView.setOnQueryTextListener(new SearchView.OnQueryTextListener() { @Override public boolean onQueryTextSubmit(String s) { - subject.onComplete(); + subject.onNext(s); return true; } diff --git a/app/src/main/java/com/rxjava2/android/samples/ui/search/SearchActivity.java b/app/src/main/java/com/rxjava2/android/samples/ui/search/SearchActivity.java index 9086b01..d9353c9 100644 --- a/app/src/main/java/com/rxjava2/android/samples/ui/search/SearchActivity.java +++ b/app/src/main/java/com/rxjava2/android/samples/ui/search/SearchActivity.java @@ -4,11 +4,12 @@ import android.widget.SearchView; import android.widget.TextView; +import androidx.appcompat.app.AppCompatActivity; + import com.rxjava2.android.samples.R; import java.util.concurrent.TimeUnit; -import androidx.appcompat.app.AppCompatActivity; import io.reactivex.Observable; import io.reactivex.ObservableSource; import io.reactivex.android.schedulers.AndroidSchedulers; @@ -56,7 +57,12 @@ public boolean test(String text) { .switchMap(new Function>() { @Override public ObservableSource apply(String query) { - return dataFromNetwork(query); + return dataFromNetwork(query) + .doOnError(throwable -> { + // handle error + }) + // continue emission in case of error also + .onErrorReturn(throwable -> ""); } }) .subscribeOn(Schedulers.io()) diff --git a/app/src/main/res/layout/activity_selection.xml b/app/src/main/res/layout/activity_selection.xml index 2c42930..999f99c 100644 --- a/app/src/main/res/layout/activity_selection.xml +++ b/app/src/main/res/layout/activity_selection.xml @@ -4,10 +4,10 @@ android:layout_width="match_parent" android:layout_height="match_parent" android:fadeScrollbars="false" - android:paddingBottom="@dimen/activity_vertical_margin" android:paddingLeft="@dimen/activity_horizontal_margin" - android:paddingRight="@dimen/activity_horizontal_margin" android:paddingTop="@dimen/activity_vertical_margin" + android:paddingRight="@dimen/activity_horizontal_margin" + android:paddingBottom="@dimen/activity_vertical_margin" tools:context="com.rxjava2.android.samples.ui.SelectionActivity"> +