【问题标题】:How to switch back from blocking scheduler to previous scheduler using netty-reactor?如何使用 netty-reactor 从阻塞调度程序切换回以前的调度程序?
【发布时间】:2022-01-02 14:53:39
【问题描述】:

如何使用 Spring Webflux + Netty + Reactor 从阻塞调度器(blocking-pool)切换回之前的调度器(reactor-http-nio)?

代码:

@RequiredArgsConstructor
@Service
@Slf4j
public class BookService {

    private final IBookRepo bookRepo;

    private final BlockingPoolConfig blockingPoolConfig;

    public Mono<Optional<Book>> getBook(Long id) {
        log.debug("getBook() - id: {}", id);
        return asyncCallable(() -> {
            log.trace("getBook() - invoking bookRepo.findById(id) ...");
            return bookRepo.findById(id);
        });
    }

    protected <S> Mono<S> asyncCallable(Callable<S> callable) {
        return Mono.fromCallable(callable)
                .subscribeOn(blockingPoolConfig.blockingScheduler()); 
    }
}

@RestController
@RequiredArgsConstructor
@Slf4j
public class BookController {

    private final BookService bookService;

    @GetMapping("/book/{id}")
    public Mono<Book> get(@PathVariable Long id) {
        log.debug("get() - id: {}", id);
        return bookService.getBook(id)
                .publishOn(Schedulers.parallel())  //publishOn(... ?)
                .map(optionalBook -> {
                    return optionalBook.map(book -> {
                        log.debug("get() result: {}", book);
                        return book;
                    }).orElseThrow(() -> {
                        log.debug("book with id: {} is not found.", id);
                        return new ResponseStatusException(HttpStatus.NOT_FOUND, "Book not found");
                    });
                });
    }

@Configuration
@Slf4j
public class BlockingPoolConfig {

    @Value("${spring.datasource.maximumPoolSize:8}")
    private int connectionPoolSize = 1;

    @Scope("singleton")
    @Bean
    public Scheduler blockingScheduler() {
        Scheduler scheduler = Schedulers.newBoundedElastic(connectionPoolSize, connectionPoolSize, "blocking-pool");
        return scheduler;
    }
}

上面我使用的是publishOn(Schedulers.parallel()),但是这个创建了新的线程池(并行)。而不是这个,我更喜欢切换 reactor-http-nio 线程池。

实际结果日志:

19:17:45.290 [reactor-http-nio-2       ] DEBUG t.a.p.controller.BookController    - get() - id: 1
19:17:45.291 [reactor-http-nio-2       ] DEBUG t.a.p.service.BookService          - getBook() - id: 1
19:17:45.316 [blocking-pool-1          ] TRACE t.a.p.service.BookService          - getBook() - invoking bookRepo.findById(id) ...
19:17:45.427 [parallel-2               ] DEBUG t.a.p.controller.BookController    - get() result: Book(id=1, title=Abc)

预期结果日志:

19:17:45.290 [reactor-http-nio-2       ] DEBUG t.a.p.controller.BookController    - get() - id: 1
19:17:45.291 [reactor-http-nio-2       ] DEBUG t.a.p.service.BookService          - getBook() - id: 1
19:17:45.316 [blocking-pool-1          ] TRACE t.a.p.service.BookService          - getBook() - invoking bookRepo.findById(id) ...
19:17:45.427 [reactor-http-nio-2       ] DEBUG t.a.p.controller.BookController    - get() result: Book(id=1, title=Abc)

【问题讨论】:

  • 什么是 BlockingPoolConfig?
  • 我添加了 BlockingPoolConfig 的来源

标签: spring-boot spring-webflux project-reactor reactor-netty


【解决方案1】:

目前这是不可能的,因为 A)这些 HTTP 线程不是由 Reactor Scheduler 控制,而是由底层的 Netty 事件循环本身控制,并且 B)Java 中没有通用的方法来“将执行返回到 (任意)线程”,如果该线程没有与之关联的Executor/ExecutorService

对于 reactor-netty,一旦您切换出 HTTP 线程,就应该没有理由想要切换回 Netty 线程。响应发送后,reactor-netty 会自然而然地完成。

假设阻塞池类似于Schedulers.boundedElastic(),您可能确实想去Schedulers.parallel() 来限制阻塞线程的寿命,这是一个非常好的解决方案。

【讨论】:

  • 据我了解,Schedulers.parallel() 不建议用于阻塞执行。 IBookRepo 是用于数据库交互的 spring 存储库。
  • 当然。但是如果您的进程的阻塞部分是有限的,并且您可以选择以非阻塞方式继续其余处理,那么切换回 parallel() 是一个不错的选择
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-09-19
  • 2023-03-06
  • 1970-01-01
  • 2019-11-21
  • 1970-01-01
  • 2020-09-12
相关资源
最近更新 更多