I'm trying to do a call in RxJava. I may be doing too many chains in the call. There are so many Observables and transformations that I am afraid the assembly is lost. Unfortunately, I need to use these functions. I have a method here that I am calling in my fragment. The method never gets called when I'm debugging and set a break point.
In my Fragment I am calling offlineItems() here:
private fun streamDownloads(): Observable<Unit> {
return downloadsDataRepository.offlineItems()
.observeOn(AndroidSchedulers.mainThread())
.map { downloadLoaded(it) } // Exception here.
}
The downloadLoaded(it) is not called and sometimes it returns an exception. I;m concerned that I map have too complicated a chained call in my RxJava. Here is the offlineItems() call.
fun offlineItems(): Observable<List<MediaItem>> {
val list = getAllMyDownloadedMediaItems()
return list.flatMapIterable { it }
.flatMap { mediaStore.getMediaItemWithId(MediaId(it.request.id)) } // this method is called .
.toList() // returns an Observable<MediaItem>
.toObservable()
}
Just to be thorough, the mediaStoreCall is:
override fun getMediaItemWithId(mediaId: MediaId): Observable<MediaItem> {
return if (isOnline()) {
Observable.just(Unit)
.effectMap { fetchAndStoreRemote(mediaId) }
.flatMap { mediaItemFromDB(mediaId) }
} else {
mediaItemFromDB(mediaId)
}
}
The list I am getting (downloads from exoplayer):
fun getAllMyDownloadedMediaItems(): Observable<List<Download>> {
return Observable.just(downloadManager.downloadIndex.getDownloads(Download.STATE_COMPLETED).use { index ->
mutableListOf<Download>().apply {
if (index.isFirst || index.moveToFirst()) {
do {
add(index.download)
} while (index.moveToNext())
}
}
})
}
Downloaded method in the fragment called from the above method:
private fun downloadLoaded(downloads: List<MediaItem>) {
if (downloads.isEmpty()) {
downloadStateView.setState(StateView.State.EMPTY)
} else {
downloadStateView.setState(StateView.State.CONTENT)
episodeAdapter.items = listOf(Header(downloads.size)) + downloads.map(::Item)
}
}
I'm calling the streamDownloads in onResume of the fragment.
override fun onResume() {
super.onResume()
Observable.merge(
streamDownloads(),
handleEmptyAction()
).autoDispose(this)
.subscribe()
}
I'm by no means an RxJava expert, so if someone could point me to what I'm doing wrong.
EDIT Stacktrace:
RxJavaAssemblyException: assembled
at dalvik.system.VMStack.getThreadStackTrace(Native Method)
at io.reactivex.Observable.map(Observable.java:9781)
at DownloadsFragment.streamDownloads(DownloadsFragment.kt:143)
at DownloadsFragment.onResume(DownloadsFragment.kt:128)
at androidx.fragment.app.Fragment.performResume(Fragment.java:2649)
at androidx.fragment.app.FragmentManagerImpl.moveToState
at FragmentManagerImpl.moveFragmentToExpectedState
at FragmentManagerImpl.moveToState(FragmentManagerImpl.java:1303)
at FragmentManagerImpl.dispatchStateChange(FragmentManagerImpl.java:2659)
at FragmentManagerImpl.dispatchResume(FragmentManagerImpl.java:2625)
at androidx.fragment.app.Fragment.performResume(Fragment.java:2658)
at FragmentManagerImpl.moveToState(FragmentManagerImpl.java:922)
at FragmentManagerImpl.moveFragmentToExpectedState
at FragmentManagerImpl.moveToState(FragmentManagerImpl.java:1303)
at FragmentManagerImpl.executeOpsTogether(FragmentManagerImpl.java:1884)
at FragmentManagerImpl.removeRedundantOperationsAndExecute
at FragmentManagerImpl.execPendingActions(FragmentManagerImpl.java:1727)
at FragmentManagerImpl$2.run(FragmentManagerImpl.java:150)
at android.os.Handler.handleCallback(Handler.java:873)
at android.os.Handler.dispatchMessage(Handler.java:99)
at android.os.Looper.loop(Looper.java:193)
at android.app.ActivityThread.main(ActivityThread.java:6669)
at RuntimeInit$MethodAndArgsCaller.run(RuntimeInit.java:493)
at com.android.internal.os.ZygoteInit.main(ZygoteInit.java:858)