如何清理使用.create(OnSubscribe)方法创建的Observable

wuj*_*jek 4 java rx-java

我有以下代码,Observable使用该Observable.create(OnSubscribe)方法创建自定义:

public class Main {

    public static void main(String[] args) {
        Subscription subscription = Observable
                .create(subscriber -> {
                    Timer timer = new Timer();
                    TimerTask task = new TimerTask() {

                        @Override
                        public void run() {
                            subscriber.onNext("tick! tack!");
                        }
                    };
                    timer.scheduleAtFixedRate(task, 0L, 1000L);
                })
                .subscribe(System.out::println);

        new Scanner(System.in).nextLine();
        System.err.println("finishing");

        subscription.unsubscribe();
    }
}
Run Code Online (Sandbox Code Playgroud)

Observable每秒使用一个计时器发出一个字符串.当用户按下回车键时,订阅将被取消.

但是,计时器仍然执行.我该如何取消计时器?我想某处必须有一个钩子,但我找不到它.

在.NET上,该create方法将返回一个IDisposable我可以作为我的实现来停止计时器.我不知道如何将它映射到RxJava,因为它的subscribe方法是void.

Sam*_*ter 6

更具说明性(和恕我直言更容易阅读)的解决方案是使用该Observable.using方法:

Observable<String> obs = Observable.using(
    // resource factory:
    () -> new Timer(),
    // observable factory:
    timer -> Observable.create(subscriber -> {
        TimerTask task = new TimerTask() {
            public void run() {
                subscriber.onNext("tick! tack!");
            }
        };
        timer.scheduleAtFixedRate(task, 0L, 1000L);
    }),
    // dispose action:
    timer -> timer.cancel()
);
Run Code Online (Sandbox Code Playgroud)

您声明了如何Timer创建依赖资源(the ),如何使用它来创建Observable,以及如何处理它,并且RxJava将负责在订阅时创建计时器并在取消订阅时将其处理.