控制Flux.generate发射值实现Reactor动态DAG顺序任务执行
问题根源
Flux.generate默认遵循Reactor的预取规则,会提前生成最多队列容量个元素放入缓冲,即便下游concatMap是顺序执行,上游依然会提前下发任务,导致任务提前触发doOnNext打印,和实际执行顺序脱节takeUntil放在concatMap下游,触发终止时会直接取消上游,不会等待当前正在执行的任务完成- 外部上下文
ctx跨线程修改存在可见性问题,可能导致生成下一个任务时读取到旧状态
正确实现方案
推荐方案:使用Mono.expand实现动态迭代
Mono.expand是专门用来实现「依赖前序结果生成下一个异步任务」场景的操作符,天然等待前一个任务完成后才触发下一个任务的生成,完全匹配你要的do{}while逻辑,实现如下:
// 启动第一个任务 Mono<OrchestrationContext> firstTask = getFirstTaskFromDAG().run(ctx); firstTask // 前序任务执行完成后,根据返回的上下文生成下一个任务 .expand(prevCtx -> { ReactiveTask<OrchestrationContext> nextTask = deriveNextStep(prevCtx.getLastExecutedStep(), prevCtx.getDecisionData()); if ("END".equals(nextTask.getName()) || !evaluateStatus(prevCtx)) { // 返回空表示终止迭代 return Mono.empty(); } return nextTask.run(prevCtx); }) .publishOn(this.factory.getSharedSchedulerPool()) .onErrorResume(throwable -> buildResponse(ctx, throwable)) .doOnCancel(() -> log.info("Task cancelled")) .doOnComplete(() -> log.info("Completed flow")) .subscribe();
若要保留Flux.generate的调整方案
需要强制关闭预取机制,保证上游每次只生成1个任务,并且调整操作符顺序:
Flux.generate( (SynchronousSink<ReactiveTask<OrchestrationContext>> synchronousSink) -> { ReactiveTask<OrchestrationContext> task = null; if (ctx.getLastExecutedStep() == null) { task = getFirstTaskFromDAG(); } else { task = deriveNextStep(ctx.getLastExecutedStep(), ctx.getDecisionData()); } if ("END".equals(task.getName())) { synchronousSink.complete(); return; } synchronousSink.next(task); }) // 指定publishOn预取数为1,每次只向上游请求1个元素 .publishOn(this.factory.getSharedSchedulerPool(), 1) .doOnNext(task -> log.info("On next => {}", task.getTaskName())) // concatMap同样指定预取数为1,不提前缓冲任务 .concatMap(reactiveTask -> reactiveTask.run(ctx), 1) .takeUntil(context -> !evaluateStatus(context, tasks)) .onErrorResume(throwable -> buildResponse(ctx, throwable)) .doOnCancel(() -> log.info("Task cancelled")) .doOnComplete(() -> log.info("Completed flow")) .subscribe();
额外注意点
- 上下文
ctx建议每次任务执行后返回新的不可变实例,避免多线程修改的可见性问题 - 任务执行逻辑里的状态修改需要做好线程安全控制,避免
deriveNextStep读取到中间状态
内容的提问来源于stack exchange,提问作者Pavan Kumar
相关产品推荐
相关产品推荐

