DirectProcessor已废弃,Reactor 3.5+及Kafka发布场景下的正确替代方案是什么?
问题解决
错误原因
- 泛型定义错误:
Supplier返回的泛型类型不符合要求,你需要向外发布的是RequestDTO类型的数据,所以返回值应该是Supplier<Flux<RequestDTO>>,而非Supplier<Flux<Sinks.Many<RequestDTO>>> - 规范层面:不建议直接将
Sinks.Many实例转换为Flux,推荐调用asFlux()方法得到可订阅流,避免外部直接拿到Sink实例进行写入操作
正确实现代码
1. 定义Sinks Bean
@Bean public Sinks.Many<RequestDTO> newSink() { // 如果你需要和原DirectProcessor一样的多播、背压缓冲能力,这个配置是对的 return Sinks.many().multicast().onBackpressureBuffer(); }
2. 定义Supplier Bean
@Autowired private Sinks.Many<RequestDTO> newSink; @Bean public Supplier<Flux<RequestDTO>> supplier(){ return () -> newSink.asFlux(); }
3. 数据写入示例
如果需要在业务代码中写入数据,直接注入Sinks.Many<RequestDTO>实例调用写入方法即可:
@Autowired private Sinks.Many<RequestDTO> newSink; public void publishData(RequestDTO data) { // 选择对应的写入策略,这里用的是无视回压直接推送,也可以根据你的需求选择emitNext的其他失败处理策略 newSink.emitNext(data, Sinks.EmitFailureHandler.FAIL_FAST); }
额外异常排查
如果修改代码后仍然抛出NoSuchMethodError异常,说明你项目中的依赖版本不兼容:需要检查spring-cloud-stream、spring-boot、spring-integration的版本匹配关系,建议使用Spring官方的依赖管理插件统一管控依赖版本,避免版本冲突。
内容的提问来源于stack exchange,提问作者Swapnil
相关产品推荐
相关产品推荐

