如何将基于回调的API转换为基于Observable的API?

ato*_*tok 6 java multithreading rx-java

我正在使用的库Message使用回调对象发出一系列对象.

interface MessageCallback {
    onMessage(Message message);
}
Run Code Online (Sandbox Code Playgroud)

使用某个libraryObject.setCallback(MessageCallback)调用添加回调,并使用非阻塞libraryObject.start()方法调用启动该进程.

创建Observable<Message>将发出这些对象的最佳方法是什么?

怎么libraryObject.start()阻止?

Yar*_*hiy 6

1.回调调用无数次

我们可以像这样将其转换为 Observable(RxJava 2 的示例):

Observable<Message> source = Observable.create(emitter -> {
        MessageCallback callback = message -> emitter.onNext(message);
        libraryObject.setCallback(callback);
        Schedulers.io().scheduleDirect(libraryObject::start);
        emitter.setCancellable(() -> libraryObject.removeCallback(callback));
    })
    .share(); // make it hot
Run Code Online (Sandbox Code Playgroud)

share使这个 observable 成为hot,即多个订阅者将共享一个订阅,即最多有一个回调注册到libraryObject

我使用io调度程序来安排start从后台线程进行的调用,因此它不会延迟第一次订阅。

2. 单条消息回调

这也是很常见的场景。假设我们有以下回调风格的异步方法:

libraryObject.requestDataAsync(Some parameters, MessageCallback callback);
Run Code Online (Sandbox Code Playgroud)

然后我们可以像这样将其转换为 Observable(RxJava 2 的示例):

Observable<Message> makeRequest(parameters) {
    return Observable.create(emitter -> {
        libraryObject.requestDataAsync(parameters, message -> {
            emitter.onNext(message);
            emitter.onComplete();
        });
    });
}
Run Code Online (Sandbox Code Playgroud)


Tom*_*řák 2

我认为你需要这样的东西(scala中给出的例子)

\n\n
import rx.lang.scala.{Observable, Subscriber}\n\ncase class Message(message: String)\n\ntrait MessageCallback {\n  def onMessage(message: Message)\n}\n\nobject LibraryObject {\n  def setCallback(callback: MessageCallback): Unit = {\n    ???\n  }\n\n  def removeCallback(callback: MessageCallback): Unit = {\n    ???\n  }\n\n  def start(): Unit = {\n    ???\n  }\n}\n\ndef messagesSource: Observable[Message] =\n  Observable((subscriber: Subscriber[Message]) \xe2\x87\x92 {\n    val callback = new MessageCallback {\n      def onMessage(message: Message) {\n        subscriber.onNext(message)\n      }\n    }\n    LibraryObject.setCallback(callback)\n    subscriber.add {\n      LibraryObject.removeCallback(callback)\n    }\n  })\n
Run Code Online (Sandbox Code Playgroud)\n\n

至于阻塞/非阻塞start():通常基于回调的架构将回调订阅和进程启动分开。在这种情况下,您可以完全独立于进程messageSource何时创建任意数量的s。start()此外,是否分叉完全由你决定。你的架构与此不同吗?

\n\n

您还应该以某种方式处理完成该过程。最好的方法是onCompleted向 MessageCallback 接口添加一个处理程序。如果你想处理错误,还可以添加一个onError处理程序。现在看,您刚刚声明了 RxJava 的基本构建基石,即观察者:-)

\n