gpt4 book ai didi

java - RxJava2 Flowable zip 异步方法调用

转载 作者:行者123 更新时间:2023-11-30 06:15:53 25 4
gpt4 key购买 nike

一个简单的例子:

import io.reactivex.Flowable;
import org.reactivestreams.Publisher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class RxTest {

private static Logger log = LoggerFactory.getLogger(RxTest.class);

public static void main(String[] args) throws InterruptedException {
log.debug("start");
RxTest rx = new RxTest();
Flowable.zip(rx.m1(), rx.m2(), rx::result).subscribe( (r) -> log.debug(r));
}

private Publisher<String> m1() throws InterruptedException {
log.debug("m2");
String result = new TestRxObj().getObj(4000L, "Hello");
return Flowable.just(result);
}

private Publisher<Long> m2() throws InterruptedException {
log.debug("m1");
Long result = new TestRxObj().getObj(100L, 777L);
return Flowable.just(result);
}

private String result(String t1, Long t2) {
log.debug("result");
return t1 + "-" + t2;
}

}

public class TestRxObj {

public <T> T getObj(Long time, T t) throws InterruptedException {
Thread.sleep(time);
return t;
}

}

所以问题是,我想异步调用方法 m1() 和 m2()

现在我有这个输出:

2018-03-09 22:34:20 [main] DEBUG RxTest: 18 - start
2018-03-09 22:34:20 [main] DEBUG RxTest: 24 - m2
2018-03-09 22:34:24 [main] DEBUG RxTest: 30 - m1
2018-03-09 22:34:24 [main] DEBUG RxTest: 36 - result
2018-03-09 22:34:24 [main] DEBUG RxTest: 20 - Hello-777

在 git 上,他们有一些示例,如何在新线程中运行:

import io.reactivex.schedulers.Schedulers;

Flowable.fromCallable(() -> {
Thread.sleep(1000); // imitate expensive computation
return "Done";
})
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.single())
.subscribe(System.out::println, Throwable::printStackTrace);

Thread.sleep(2000); // <--- wait for the flow to finish

但是这个例子的问题是,我不知道如何加入(主要/当前)威胁,而不是编写Thread.sleep(2000);。我不是在 Android 上运行这个,只是简单的 java 项目,所以如果我正确理解它在做什么,我就不能 .observeOn(AndroidSchedulers.mainThread())

我的主要目标,使用 RxJava2 异步运行 1..N(至少 2 :)) 方法,如果不用 zip 就可以完成,我可以接受,任何使用 RxJava 的方法,无需 Thread.sleep(...) :)

最佳答案

也许这会有所帮助:

public static void main(String[] args) throws InterruptedException     {
log.debug("start");
RxTest rx = new RxTest();
Flowable.zip(
Observable.defer(() -> rx.m1()).subscribeOn(Schedulers.io()),
Observable.defer(() -> rx.m2()).subscribeOn(Schedulers.io()),
rx::result)
.subscribe( (r) -> log.debug(r));
}

出于测试目的,有时将 Schedulers.io() 作为参数或字段很有用,因此您可以将其替换为 TestScheduler

关于java - RxJava2 Flowable zip 异步方法调用,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/49202012/

25 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com