【发布时间】:2017-10-06 22:12:46
【问题描述】:
我有一些遗留的非 RX 代码,它们通过产生一个新线程来完成一些网络工作。 当工作完成时,它会调用一个回调方法。
我无法控制运行此代码的线程。它是遗留的,它自己会产生一个新的Thread。
这可以简化为:
interface Callback {
void onSuccess();
}
static void executeRequest(String name, Callback callback) {
new Thread(() -> {
try {
System.out.println(" Starting... " + name);
Thread.sleep(2000);
System.out.println(" Finishing... " + name);
callback.onSuccess();
} catch (InterruptedException ignored) {}
}).start();
}
我想将其转换为 RX Completable。为此,我使用Completable#create()。
CompletableEmitter 的实现调用executeRequest 的实现传递Callback 的实现,该信号在请求完成时发出信号。
订阅时我还会打印日志跟踪以帮助我调试。
static Completable createRequestCompletable(String name) {
return Completable.create(e -> executeRequest(name, e::onComplete))
.doOnSubscribe(d -> System.out.println("Subscribed to " + name));
}
这按预期工作。 Completable 仅在“请求”完成并调用回调后完成。
问题是,在trampoline 调度程序中订阅这些可完成项时,它不会在订阅第二个请求之前等待第一个请求完成。
这段代码:
final Completable c1 = createRequestCompletable("1");
c1.subscribeOn(Schedulers.trampoline()).subscribe();
final Completable c2 = createRequestCompletable("2");
c2.subscribeOn(Schedulers.trampoline()).subscribe();
输出:
Subscribed to 1
Starting... 1
Subscribed to 2
Starting... 2
Finishing... 1
Finishing... 2
如您所见,它在第一个 Completable 完成之前订阅了第二个 Completable,即使我正在订阅 trampoline。
我想将completables 排队,以便第二个等待第一个完成,输出以下内容:
Subscribed to 1
Starting... 1
Finishing... 1
Subscribed to 2
Starting... 2
Finishing... 2
我确定问题与工作线程中正在完成的工作有关。如果Completable 的实现没有产生新线程,它会按预期工作。
但这是遗留代码,我想做的是在不修改的情况下使其适应 RX。
注意:请求在程序的不同点执行 - 我不能使用 andThen 或 concat 来实现序列化执行。
【问题讨论】:
标签: java multithreading rx-java scheduler rx-android