【问题标题】:Flux blocking Netty start-up助焊剂阻止 Netty 启动
【发布时间】:2019-03-14 07:31:53
【问题描述】:

我有一个从实时数据源读取的@Repository。我正在使用 Flux.create() { sink->sink.next() }

提供数据

@Service 正在执行以下操作;

@Autowired MyRepository myRepository;

@PostConstruct() public void startUp() {
  ConnectableFlux<Object> cf = myRepository.flux.publish();
  cf.subscribe(System.out::println);
  cf.connect();
}

这可以工作并打印数据,但我确实没有在日志中得到“Netty started”并且@Controllers 没有响应。如果我省略cf.connect(),Netty 就会启动。所以我假设cf.connect() 正在阻止 Netty。

理想情况下,我希望订阅自动启动。在@PostConstuct 中使用connect() 是否为时过早?我应该收听“Netty Started”事件,然后connect(), 还是我的订阅完全错误?

编辑:如果connect 在一个守护进程Thread 中运行,Netty 会启动并且订阅有效。

【问题讨论】:

  • 删除@PostConstruct 并使用SmartInitializingSingleton 调用startUp 没有帮助。

标签: java spring-webflux project-reactor


【解决方案1】:

@EnableAsync 放在Spring Boot 的主应用程序类上,将@Async 放在上述方法上似乎可行。

编辑: 我在这里找到了更好的解决方案Connectable Flux blocks on toIterator.forEach, while Flux does not. #1549。所以我的代码现在看起来像这样;

Flux<MyClass> flux = Flux.create(
   sink -> {
       while(condition) {
          sink.next(nextValueFromDataSource);
       }
       sink.complete();
    }
)
.publish()
.autoConnect(1);

也因为autoConnect(1),所以不需要@Async

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-02-29
    • 2015-03-22
    • 1970-01-01
    • 1970-01-01
    • 2016-08-12
    • 2016-03-12
    • 2015-02-11
    • 1970-01-01
    相关资源
    最近更新 更多