【问题标题】:Code adapted to RX Completable not blocking onSubscribe thread适应 RX Completable 的代码不阻塞 onSubscribe 线程
【发布时间】: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。

注意:请求在程序的不同点执行 - 我不能使用 andThenconcat 来实现序列化执行。

【问题讨论】:

    标签: java multithreading rx-java scheduler rx-android


    【解决方案1】:

    我已经设法通过使用Latch 显式阻止订阅Thread 来按顺序执行Completables。 但我不认为这是在 RX 中执行此操作的惯用方式,我仍然不明白为什么我需要这样做,并且线程在 Completable 完成之前不会被阻塞。

    static Completable createRequestCompletable(String name) {
        final CountDownLatch latch = new CountDownLatch(1);
        return Completable.create(e -> {
            executeRequest(name, () -> {
                e.onComplete();
                latch.countDown();
            });
            latch.await();
        })
        .doOnSubscribe(disposable -> System.out.println("Subscribed to " + name));
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-07-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-01-18
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多