【问题标题】:How to Merge multiple Observables with diferents call backs into one single Stream?如何将多个具有不同回调的 Observable 合并到一个 Stream 中?
【发布时间】:2018-04-11 18:22:43
【问题描述】:

我有这个问题。我正在尝试将本地数据库与远程应用程序同步到我的 android 应用程序中。我正在创建上传本地创建的新信息的逻辑,服务器使用远程 ID 响应并将其保存在服务器中。为了存档这个,我使用了一个方法,它接受一个对象数组并返回一个 Observable,它为每个元素发出服务器的响应。 像这样。

val dailyEntries = App.db.dailyEntryDao().getDailyEntriesCreated()
            dailyEntries.sync(context) //Return an observable
                    .observeOn(AndroidSchedulers.mainThread())
                    .subscribe({                        
                        val response = DailyEntry(it)//Creates a daily entry using 
                                                      the response from the server
                        thread {                                
                          App.db.dailyEntryDao().update(response)
                        }
                    }, {
                        it.printStackTrace()
                    }, {
                        uploadEnclosures()
                    })

你怎么看,在onSuccess中从当前observable调用了另一个方法。它使用相同的逻辑并且显示在前面。

private fun uploadEnclosures() {
        thread {
            val enclosures = App.db.enclosureDao().getEnclosuresCreated()
            enclosures.sync(context)
                    .observeOn(AndroidSchedulers.mainThread())
                    .subscribe({
                        val response = Enclosure(it)
                        thread {                                
                            App.db.enclosureDao().update(response)
                        }
                    }, {
                        it.printStackTrace()
                    }, {
                        uploadSmokeTest()
                    })
        }
    }

所有表格都继续。我们总是在当前 Observable 的 onSuccess 中执行下一个表的更新。这样做是因为我需要按特定顺序进行同步。

我的问题是,有没有办法将所有这些 Observable 合并为一个来执行单个订阅并控制每个 onNext 情绪?

感谢您的回答

【问题讨论】:

  • thread {?在 RxJava 订阅者中启动后台线程?那么observeOn()subscribeOn() 是干什么用的?
  • 这是因为 val enclosures = App.... 从房间表中检索信息。观察在 MainThread 中很有用,以防我想通知 de UI 任何更改。您是否建议进行任何更改?

标签: android kotlin rx-java reactive-programming rx-kotlin


【解决方案1】:

是的,但是需要做一些工作,您可以使用concat 运算符为您处理排序,并按顺序传递observables 的列表,然后使用订阅它单个观察者期望 Any 事件会向下传播。

为了严格控制类型安全,您可以使用通用接口标记可观察源类型,并使用实例检查来执行特定于事件类型的操作。

查看更多here

代码示例-

fun concatCalls(): Observable<Any> {
    return Observable.concat(src1, src2, ...)
}

那么消费者应该是这样的 -

concatCalls().subscribe(object: Subscriber<Any> {
    override fun onNext(o: Any) {
       when (o) {
           is Object1 -> // do handling for stuff emitted by src1
           is Object2 -> // do handling for stuff emitted by src2
           ....
           else // skip
    }
    ....
})

【讨论】:

  • 老实说我不明白怎么做xD,你能发一个例子吗?另外,我想知道我可以在每个 Observable 中实现 doOnNext 方法,这样每个方法都知道回调的类型。你怎么看?
  • Tnx 的答案,我要实现它,看看会发生什么。 :D
  • 嘿@JhonFredyTrujilloOrtega 如果这有帮助,你能接受答案吗?
猜你喜欢
  • 2019-06-06
  • 2011-10-29
  • 2020-05-01
  • 1970-01-01
  • 2017-08-04
  • 2020-07-11
  • 1970-01-01
  • 1970-01-01
  • 2018-11-16
相关资源
最近更新 更多