【问题标题】:How to create a job that runs at a fixed interval in RxJava?如何在 RxJava 中创建以固定间隔运行的作业?
【发布时间】:2015-06-05 17:01:55
【问题描述】:

我正在尝试创建一个任务,该任务将定期查询我的数据库,将所有结果写入其他状态,我想使用 RxJava 来执行此操作。

我正在使用RxJava-JDBC 来查询我的数据库。代码如下所示:

    final Database db = Database.from(url);

    db
        .select("SELECT f1,f2 FROM mydata")
        .autoMap(MyDatum.class)
        .subscribe(
            new Action1<MyDatum>() {
                @Override
                public void call(MyDatum t) {
                    state.add(t);
                }
            },
            new Action1<Throwable>() {
                @Override
                public void call(Throwable t) {
                    L.error("Task failed", t);
                }
            },
            new Action0() {
                @Override
                public void call() {
                    state.makeAvailable();
                }
            }
        );

问题是,当我订阅时它会起作用,然后它就会停止。所以我使用了Observable.interval 并且有这个工作:

    Observable
        .interval(10, TimeUnit.SECONDS)
        .forEach(
            new Action1<Long>() {
                @Override
                public void call(Long arg) {
                    db
                        .select("SELECT f1,f2 FROM mydata")
                        .autoMap(MyDatum.class)
                        .subscribe(
                            new Action1<MyDatum>() {
                                @Override
                                public void call(MyDatum t) {
                                    state.add(t);
                                }
                            },
                            new Action1<Throwable>() {
                                @Override
                                public void call(Throwable t) {
                                    L.error("Task failed", t);
                                }
                            },
                            new Action0() {
                                @Override
                                public void call() {
                                    state.makeAvailable();
                                }
                            }
                        );
                }
            }
        );

但我想知道如果我没有做错什么让一个流嵌套在另一个流中。 我考虑过使用flatMap,但是onComplete 永远不会被执行,因为interval 永远不会调用onComplete

我希望它不仅可以通过间隔触发,还可以通过传入事件触发。

我在这里遗漏了什么吗? 谢谢

【问题讨论】:

  • doOnNextdoOnCompleted 运算符在您的情况下非常有用。这是一个示例,您可以如何使用这些运算符实现所描述的行为 gist.github.com/nsk-mironov/087bbe570a617d366260
  • 您使用的 interval 方法有一个怪癖。它只会在 x timeunit 过去后启动。我建议立即使用timer(0, 10, TimeUnit.SECONDS)。见Github PR
  • @VladimirMironov 谢谢,这看起来确实比我做的要好一些。你介意提出这个作为答案吗?

标签: reactive-programming rx-java


【解决方案1】:

doOnNextdoOnCompleted 运算符在您的情况下非常有用。下面是一个示例,您可以如何使用这些运算符来实现所描述的行为:

final Observable<MyDatum> observable = Observable.interval(10, TimeUnit.SECONDS).flatMap(new Func1<Long, Observable<MyDatum>>() {
    @Override
    public Observable<MyDatum> call(final Long counter) {
        return db.select("SELECT f1,f2 FROM mydata")
                .autoMap(MyDatum.class)
                .doOnNext(new Action1<MyDatum>() {
                    @Override
                    public void call(final MyDatum value) {
                        state.add(value);
                    }
                })
                .doOnCompleted(new Action0() {
                    @Override
                    public void call() {
                        state.makeAvailable();
                    }
                });
    }
});

final Subscription subscription = observable.subscribe();

【讨论】:

    【解决方案2】:

    flatMapmerge 是您要使用的运算符。首先,您应该避免订阅操作员体内的可观察对象。而是使用 flatMap 并返回 observable。这将为您订阅所有发出的 observables。

    为了手动触发查询,您可以合并 PublishSubject (documentation),您可以调用 onNext 来推送事件并手动触发查询。把你的代码改成这样。

    PublishSubject<Long> subject = PublishSubject.create();
    Observable.merge(subject, Observable.timer(0, 1, TimeUnit.SECONDS))
        .flatMap(new Func1<Long, Observable<MyDatum>>() {
            @Override
            public Observable<MyDatum> call(Long arg) {
                return db
                    .select("SELECT f1,f2 FROM mydata")
                    .autoMap(MyDatum.class);
        }}).subscribe(
            new Action1<MyDatum>() {
                @Override
                public void call(MyDatum t) {
                    state.add(t);
                }
            },
            new Action1<Throwable>() {
                @Override
                public void call(Throwable t) {
                    L.error("Task failed", t);
                }
            },
            new Action0() {
                @Override
                public void call() {
                    state.makeAvailable();
                }
            }
        );
    // you can call onNext with any value to trigger a manual query
    subject.onNext(999L);
    

    这是一个演示此行为的简单 RxJava sn-p。

    CountDownLatch l = new CountDownLatch(5);
    PublishSubject<Long> subject = PublishSubject.create();
    Observable.merge(subject, Observable.timer(0, 1, TimeUnit.SECONDS).take(3))
            .flatMap((Long arg) -> {
                System.out.println("tick: " + arg);
                l.countDown();
                return Observable.just(arg+10);
            })
            .forEach(System.out::println);
    l.await(1, TimeUnit.SECONDS);
    subject.onNext(999L);
    l.await();
    

    输出

    tick: 0
    10
    tick: 999
    1009
    tick: 1
    11
    tick: 2
    12
    

    【讨论】:

    • 谢谢,我不知道 PublishSubject。它可能会让我得到我一直在寻找的外部触发。不过,关于问题的主要部分,我没有使用 flatMap 的原因是,正如我所说,完整的动作永远不会被调用,因为计时器永远不会完成。我想我想是一个在每个 onNext 之后发送 onComplete 的计时器。
    猜你喜欢
    • 2013-03-25
    • 2014-10-18
    • 2021-02-24
    • 1970-01-01
    • 2020-07-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多