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()阻止?
我们可以像这样将其转换为 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从后台线程进行的调用,因此它不会延迟第一次订阅。
这也是很常见的场景。假设我们有以下回调风格的异步方法:
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)
我认为你需要这样的东西(scala中给出的例子)
\n\nimport 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 })\nRun Code Online (Sandbox Code Playgroud)\n\n至于阻塞/非阻塞start():通常基于回调的架构将回调订阅和进程启动分开。在这种情况下,您可以完全独立于进程messageSource何时创建任意数量的s。start()此外,是否分叉完全由你决定。你的架构与此不同吗?
您还应该以某种方式处理完成该过程。最好的方法是onCompleted向 MessageCallback 接口添加一个处理程序。如果你想处理错误,还可以添加一个onError处理程序。现在看,您刚刚声明了 RxJava 的基本构建基石,即观察者:-)
| 归档时间: |
|
| 查看次数: |
1871 次 |
| 最近记录: |