【问题标题】:Control the delay in execution of a repeat task for Mono控制执行 Mono 的重复任务的延迟
【发布时间】:2020-04-15 16:35:32
【问题描述】:

我正在尝试实现一种轮询机制。我想根据某些条件增加或减少轮询间隔。我正在使用 Mono.repeat 和 delayElements 来执行间隔的重复任务。但我无法找到一种方法来根据某些标准修改延迟.

Mono.just(1).
    repeat().
    delayElements(getPollingInterval()).
    takeUntil((s)->
      {

          if(checkForEndCriteria()){
              log.info("Critera to end reached);
              return true;
          }
          return false;
      }).
    log().
    subscribeOn(Schedulers.boundedElastic()).
    flatMapSequential(x -> {
        List<Event> eventList = getEvents(id, lastItemTimeStamp);;
        if (!eventList.isEmpty()) {
            //Recieving events now. So want to decrease the interval.
            return Flux.fromIterable(eventList);
        } else {
        //There are no events happening .So I would like
        //to increase the delay of repeat task by 1 sec

            return Flux.just(buildHeartBeatEvent());
        }
    }).
    onErrorResume(error -> {
        log.error("Error occurred", error);
        return Flux.error(error);
    });```

【问题讨论】:

  • 为什么要用reactor来轮询?

标签: spring-boot project-reactor reactor


【解决方案1】:

我用Flux.&lt;Duration&gt;generate() 实现了这个:

        Flux
            .<Duration>generate(sink -> {
                Date date = [NEXT_DATE];
                if (date != null) {
                    long millis = date.getTime() - System.currentTimeMillis();
                    sink.next(Duration.ofMillis(millis));
                }
                else {
                    sink.complete();
                }
            })
            .concatMap(duration ->
                    Mono.delay(duration)
                    ...
            )
            .repeat();

所以,每次我们带着那个repeat()回到generate()时,我们都可以查看一些状态以获得下一次执行Date

【讨论】:

  • 感谢您的建议。我试过这个。但我面临的问题是,在这种方法中,首先生成是创建一个通量,直到满足条件。在我的用例中,结束条件不是固定的。它基于用户活动。我不断地轮询数据库。我正在尝试根据新事件的可用性来增加轮询的延迟。如果没有事件发生,则增加延迟。如果发生新事件,则将轮询间隔重置为初始值。
【解决方案2】:

我建议您在 Flux 之后使用 concatWith 在 flatMapSequential 中使用 Mono.delay,而不是使用 delayElements 来实现固定延迟,这样您就可以根据您的元素轻松控制其持续时间。

您也可以使用repeat(BooleanSupplier) 代替repeat + takeUntil 和defaultIfEmpty 来清理一点代码。

我希望这样的事情可以帮助你:

Mono.fromCallable(() -> getEvents(id, lastItemTimeStamp)).
subscribeOn(Schedulers.boundedElastic()).
flatMapMany(eventList -> Flux.fromIterable(eventList).
                              defaultIfEmpty(buildHeartBeatEvent()).
                              concatWith(Mono.delay(eventList.isEmpty() ?
                                          getEmptyListPollingInterval() : 
                                          getPollingInterval()).
                                         then().cast(Event.class))).
log().
repeat(this::checkForEndCriteria);

【讨论】:

    猜你喜欢
    • 2015-09-14
    • 2013-11-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-03-04
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多