【问题标题】:Controlling emit values from Flux generate for sequential task execution控制来自 Flux generate 的发射值以用于顺序任务执行
【发布时间】:2022-01-08 12:30:05
【问题描述】:

我正在使用 reactor 框架编写一个简单的编排框架,该框架按顺序执行任务,下一个要执行的任务取决于先前任务的结果。根据先前任务的结果,我可能有多种路径可供选择。早些时候,我基于静态 DAG 编写了一个类似的框架,其中我将任务列表作为可迭代对象传递并使用 Flux.fromIterable(taskList)。但是,由于静态数组发布者,这并没有给我动态的灵活性。

我正在寻找像 do(){}while(condition) 这样的替代方法来解决 DAG 遍历和任务决策,我想出了 Flux.generate()。我评估生成方法的下一步并将下一个任务传递到下游。我现在面临的问题是,Flux.generate 不会等待下游完成,而是推送直到条件设置为无效。到任务 1 执行时,任务 2 将被推送n 次,这不是预期的行为。

有人可以指点我正确的方向吗?

谢谢。

使用任务列表(静态 DAG)的第一次迭代

Flux.fromIterable(taskList)
        .publishOn(this.factory.getSharedSchedulerPool())
        .concatMap(
            reactiveTask -> {
              log.info("Running task =>{}", reactiveTask.getTaskName());
              return reactiveTask
                  .run(ctx);
            })
        // Evaluates status from previous task and terminates stream or continues.
        .takeWhile(context -> evaluateStatus(context))
        .onErrorResume(throwable -> buildResponse(ctx, throwable))
        .doOnCancel(() -> log.info("Task cancelled"))
        .doOnComplete(() -> log.info("Completed flow"))
        .subscribe();

尝试动态dag

Flux.generate(
            (SynchronousSink<ReactiveTask<OrchestrationContext>> synchronousSink) -> {
              ReactiveTask<OrchestrationContext> task = null;
              if (ctx.getLastExecutedStep() == null) {
                // first task;
                task = getFirstTaskFromDAG();
              } else {
                task = deriveNextStep(ctx.getLastExecutedStep(), ctx.getDecisionData());
                   
              }
              if (task.getName.equals("END")) {
                synchronousSink.complete();
              }
              synchronousSink.next(task);
            })
        .publishOn(this.factory.getSharedSchedulerPool())
        .doOnNext(orchestrationContextReactiveTask -> log.info("On next => {}", 
          orchestrationContextReactiveTask.getTaskName()))
        .concatMap(
            reactiveTask -> {
              log.info("Running task =>{}", reactiveTask.getTaskName());
              return reactiveTask
                  .run(ctx);                  
            })
        .onErrorResume(throwable -> buildResponse(ctx, throwable))
        .takeUntil(context -> evaluateStatus(context, tasks))
        .doOnCancel(() -> log.info("Task cancelled"))
        .doOnComplete(() -> log.info("Completed flow")).subscribe();

上述方法的问题是,在执行任务 1 时,onNext() 订阅者打印了很多次,因为generate 正在发布。我希望生成方法等待上一个任务的结果并提交新任务。在非响应式世界中,这可以通过简单的 while() 循环来实现。

每个任务都会执行以下操作。

public class ResponseTask extends AbstractBaseTask {
  private TaskDefinition taskDefinition;
  final String taskName;


  public ResponseTask(
      StateManager stateManager,
      ThreadFactory factory,
    ) {    
    this.taskDefinition = taskDefinition;
    this.taskName = taskName;
  }

   public Mono<String> transform(OrchestrationContext context) {
    Any masterPayload = Any.wrap(context.getIngestionPayload());
    return Mono.fromCallable(() -> stateManager.doTransformation(context, masterPayload);
  }

  
  public Mono<OrchestrationContext> execute(OrchestrationContext context, String payload) {
    log.info("Executing sleep for task=>{}", context.getLastExecutedStep());
    return Mono.delay(Duration.ofSeconds(1), factory.getSharedSchedulerPool())
        .then(Mono.just(context));
  }

  public Mono<OrchestrationContext> run(OrchestrationContext context) {
log.info("Executing task:{}. Last executed:{}", taskName, context.getLastExecutedStep());
  return transform(context)
         .doOnNext((result) -> log.info("Transformation complete for task=?{}", taskName);)
         .flatMap(payload -> {
             return execute(context, payload);
         }).onErrorResume(throwable -> {
             context.setStatus(FAILED);
             return Mono.just(context);
         }); 
}

}

编辑 - 来自 @Ikatiforis 的推荐 - 我得到以下输出

Here's the output from my side.

2021-12-02 09:58:14,643 INFO  (reactive_shared_pool) [ReactiveEngine lambda$doOrchestration$5:98] On next => Task1 
2021-12-02 09:58:14,644 INFO  (reactive_shared_pool) [ReactiveEngine lambda$doOrchestration$6:101] Running task =>Task1 
2021-12-02 09:58:14,644 INFO  (reactive_shared_pool) [AbstractBaseTask run:75] Executing task:Task1. Last executed:Task1 
2021-12-02 09:58:14,658 INFO  (reactive_shared_pool) [ReactiveEngine lambda$doOrchestration$5:98] On next => Task2 
2021-12-02 09:58:14,659 INFO  (reactive_shared_pool) [AbstractBaseTask lambda$run$0:83] Transformation complete for task=?Task1 
2021-12-02 09:58:14,659 INFO  (reactive_shared_pool) [ResponseTask execute:41] Executing sleep for task=>Task1 
2021-12-02 09:58:15,661 INFO  (reactive_shared_pool) [AbstractBaseTask lambda$run$4:106] Success for task=>Task1 
2021-12-02 09:58:15,663 INFO  (reactive_shared_pool) 
[ReactiveEngine lambda$doOrchestration$6:101] Running task =>Task2 
2021-12-02 09:58:15,811 INFO  (cassandra-nio-worker-8) [AbstractBaseTask run:75] Executing task:Task2. Last executed:Task2 
2021-12-02 09:58:15,811 INFO  (reactive_shared_pool) [ReactiveEngine lambda$doOrchestration$5:98] On next => Task2 
2021-12-02 09:58:15,812 INFO  (reactive_shared_pool) [AbstractBaseTask lambda$run$0:83] Transformation complete for task=?Task2 
2021-12-02 09:58:15,812 INFO  (reactive_shared_pool) [ResponseTask execute:41] Executing sleep for task=>Task2 
2021-12-02 09:58:15,837 INFO  (centaurus_reactive_shared_pool) [ReactiveEngine lambda$doOrchestration$9:113] Completed flow 

我在这里看到了几个问题--

The sequence of execution is 
 1. Task does transformations ( runs on Mono.fromCallable)
 2. Task induces a delay - Mono.fromDelay()
 3. Task completes execution. After this, generate method should evaluate the context and pass on the next task to be executed.

What I see from the output is:
 1. Task 1 starts the transformations - Runs on Mono.fromCallable.
 2. Task 2 doOnNext is reported - which means the stream already got this task.
 3. Task 1 completes.
 4. Task 2 starts and executes delay -> the stream does not wait for response from task 2 but completes the flow.

【问题讨论】:

  • 上下文中的状态是否在任务完成之前更新?这可以解释为什么最后一项任务似乎被缩短了(也许是因为takeUntil?)。如果您需要处理所有任务及其结果,为什么要使用takeUntil?看起来经典的完成就足够了......
  • @SimonBaslé 状态会根据任务错误或逻辑故障进行更新。所以它应该等待所述任务的完成,因为它是链式的。 takeUntil 接受一个方法,该方法基本上查看这个状态标志,一个二进制标志,并终止流程。如果逻辑输出为假,我们不需要继续。所以经典补全不是一种选择。
  • @SimonBaslé 同样在这种情况下,我将代码修改为 Ikatiforis 的建议 - 预取使流执行相同的任务两次,而我的预期行为是 - 生成一个任务 -> 等待直到它完成 -> 评估状态和下一步 -> 传递任务。 (这是 Flux.generate 的功能)

标签: java reactive-programming spring-webflux project-reactor


【解决方案1】:

上述方法的问题是,在执行任务 1 时, onNext() 订阅者打印了很多次,因为 generate 正在发布。

发生这种情况是因为concatMap 预先请求了许多项目(默认为 32),而不是一个一个地请求元素。如果您确实需要一次请求一个元素,您可以使用concatMap(Function&lt;? super T,? extends Publisher&lt;? extends V&gt;&gt; mapper,int prefetch) 变体方法并提供prefetch 值,如下所示:

.concatMap(reactiveTask -> {
              log.info("Running task =>{}", reactiveTask.getTaskName());
              return reactiveTask.run(ctx);                  
            }, 1)

编辑

还有一个publishOn 方法,它采用prefetch 值。看看下面的斐波那契生成器示例,如果它按预期工作,请告诉我:

generateFibonacci(100)
    .publishOn(boundedElasticScheduler, 1)
    .doOnNext(number -> log.info("On next => {}", number))
    .concatMap(number -> {
      log.info("Running task => {}", number);
      return task(number).doOnNext(num -> log.info("Task completed => {}", num));
    }, 1)
    .takeWhile(context -> context < 3)
    .subscribe();
  public Flux<Integer> generateFibonacci(int limit) {
    return Flux.generate(
        () -> new FibonacciState(0, 1),
        (state, sink) -> {
          log.info("Generating number: " + state);
          sink.next(state.getFormer());
          if (state.getLatter() > limit) {
            sink.complete();
          }
          int temp = state.getFormer();
          state.setFormer(state.getLatter());
          state.setLatter(temp + state.getLatter());
          return state;
        });
  }

这是输出:

2021-12-02 10:47:51,990  INFO main c.u.p.p.s.c.Test - Generating number: FibonacciState(former=0, latter=1)
2021-12-02 10:47:51,993  INFO pool-1-thread-1 c.u.p.p.s.c.Test - On next => 0
2021-12-02 10:47:51,996  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Running task => 0
2021-12-02 10:47:54,035  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Task completed => 0
2021-12-02 10:47:54,035  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Generating number: FibonacciState(former=1, latter=1)
2021-12-02 10:47:54,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - On next => 1
2021-12-02 10:47:54,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Running task => 1
2021-12-02 10:47:56,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Task completed => 1
2021-12-02 10:47:56,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Generating number: FibonacciState(former=1, latter=2)
2021-12-02 10:47:56,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - On next => 1
2021-12-02 10:47:56,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Running task => 1
2021-12-02 10:47:58,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Task completed => 1
2021-12-02 10:47:58,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Generating number: FibonacciState(former=2, latter=3)
2021-12-02 10:47:58,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - On next => 2
2021-12-02 10:47:58,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Running task => 2
2021-12-02 10:48:00,036  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Task completed => 2
2021-12-02 10:48:00,037  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Generating number: FibonacciState(former=3, latter=5)
2021-12-02 10:48:00,037  INFO pool-1-thread-1 c.u.p.p.s.c.Test - On next => 3
2021-12-02 10:48:00,037  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Running task => 3
2021-12-02 10:48:02,037  INFO pool-1-thread-1 c.u.p.p.s.c.Test - Task completed => 3
2021-12-02 10:52:07,877  INFO pool-1-thread-2 c.u.p.p.s.c.Test - Completed flow

编辑 04122021

你说:

我正在尝试模拟 HTTP / 阻塞调用。因此是 Mono.delay。

Mono#Delay 不是模拟阻塞调用的合适方法。延迟是通过并行调度程序引入的,因此它不会等待任务完成。您可以像这样模拟阻塞调用:

  public String get() throws IOException {
    HttpsURLConnection connection = (HttpsURLConnection) new URL("https://jsonplaceholder.typicode.com/comments").openConnection();
    connection.setRequestMethod("GET");
    try(InputStream inputStream = connection.getInputStream()) {
      return new String(inputStream.readAllBytes(), StandardCharsets.UTF_8);
    }
  }

请注意,作为替代方案,您可以使用 .limitRate(1) 运算符而不是 prefetch 参数。

【讨论】:

  • @Ikatiforis 我添加了这个更改,但现在,在任务 2 完成执行之前,通量会发出onComplete 的信号。同样,执行哪个任务的决定应该在任务 1 执行之后再决定,但是在这里,任务 2 已经预取并且准备好了。如何在Flux.generate 中等待、重新评估和生成正确的任务?
  • @PavanKumar 编辑完成
  • @ikatiforis 我用我的输出更新了主线程,我的行为与你的不同。我有一些Mono.fromCallable 和延迟通话。
  • @PavanKumar 所以,更新您的答案以包含您的所有功能。
  • @ikatiforis 也更新了任务信息。
猜你喜欢
  • 2011-01-10
  • 1970-01-01
  • 1970-01-01
  • 2021-03-12
  • 2021-12-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-03-07
相关资源
最近更新 更多