【问题标题】:Subscribe on Observable Does Nothing订阅 Observable 什么都不做
【发布时间】:2019-07-02 05:38:37
【问题描述】:

我试图了解 Observables 是如何执行的,但似乎无法让这个简单的代码工作。

public class RxJavaExample {
    public static void main(String[] args) {
        Observable<String> hello = Observable.fromCallable(() -> 
            getHello()).subscribeOn(Schedulers.newThread());

        hello.subscribe();

        System.out.println("End of main!");
    }

    public static String getHello() {
        System.out.println("Hello called in " + 
            Thread.currentThread().getName());
        return "Hello";
    }
}

不应该hello.subscribe()执行getHello()吗?

【问题讨论】:

  • getHello() 没有被执行?
  • @SantanuSur 不,它从未调用过。

标签: java observable rx-java2


【解决方案1】:

这是因为您的主线程在后台线程到达getHello 之前完成。尝试在您的main 方法中添加Thread.sleep(5000),然后退出。

或者,等到您订阅的onCompleted 被调用。

编辑:程序终止的原因是因为 RxJava 产生了daemon 线程。在寻找好的来源时,我还发现了this问题,它可能也回答了它。

【讨论】:

  • 我认为 Java 程序在所有线程完成之前不会退出。此外,我确实在getHello() 方法中设置了一个断点,它从未被调用过。
  • 您为什么不尝试一下,而不是假设您是正确的并且知道它是如何工作的?请看看我的编辑。也许也看看stackoverflow.com/questions/2213340/…
  • Java 无需等待线程完成即可愉快地退出,这就是我们有 ExecutorService.awaitTermination 的原因。
  • LeffeBrune,在@sfiss 指出之前,我最初并不知道 RxJava 创建了守护线程。谢谢你们俩。
【解决方案2】:

@sfiss 是对的,这正如你所期望的那样工作:

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

import io.reactivex.Observable;
import io.reactivex.schedulers.Schedulers;

public class RxJavaExample {
  public static void main(String[] args) throws InterruptedException {
    ExecutorService exec = Executors.newCachedThreadPool();
    Observable<String> hello = Observable.fromCallable(() -> getHello())
        .subscribeOn(Schedulers.from(exec));

    hello.subscribe();

    System.out.println("End of main!");

    exec.shutdown();
    exec.awaitTermination(10, TimeUnit.SECONDS);
  }

  public static String getHello() {
    System.out.println("Hello called in " + Thread.currentThread().getName());
    return "Hello";
  }
}

输出如下:

End of main!
Hello called in pool-1-thread-1

【讨论】:

  • 通过@sfiss 接受这个答案,因为这个答案更完整。
  • 这不包括仅线程的可观察对象,因此实际上并没有回答标题问题。似乎与 Threads 和 observables 混淆
  • @TheresaForster 虽然您在技术上是正确的,但 OP 遇到的问题是由于订阅者的异步性质造成的。
  • 看起来我们也可以进行阻塞订阅——subscribe(Observable::blockingSubscribe)
【解决方案3】:

你可能对线程和 Observable 感到困惑,

我过去使用 Observables 的方式是用于 Minecraft 插件上的计时器,我有一个每分钟触发的事件。

public class TimerHandler extends Observable implements Runnable{

    @Override
    public void run() {
        this.setChanged();
        this.notifyObservers();
    }
}

所以这每分钟触发一次,然后将事件添加到计时器队列中,您只需订阅 observable 意味着订阅的调用每分钟触发一次。

public class PlotTimer implements Observer {

    @Override
    public void update(Observable o, Object arg) {
        ......

要订阅我打电话给以下

getServer().getScheduler().scheduleAsyncRepeatingTask(this,timerHandler,1200,1200);
timerHandler.addObserver(new PayDayTimer());
timerHandler.addObserver(new ProfileTimer());
timerHandler.addObserver(new PlotTimer());

【讨论】:

  • 问题不是关于java.util.Observable,而是关于RxJava。我的和接受的答案都详细说明了 RxJava 如何使用守护线程以及如何防止提前关闭。
  • @Theresa AFAIK 推荐的方法是使用标签来指定技术,而不是在标题中提及。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-12-13
  • 1970-01-01
  • 2019-12-29
  • 1970-01-01
  • 2020-10-06
  • 2018-10-26
相关资源
最近更新 更多