Spring Reactive中动态增减生产者且不中断流的实现可行性咨询
实现可行性
该需求完全可以实现,核心思路是将固定的Flux.merge参数列表,替换为一个可动态写入新Publisher的上层流,通过流的平铺逻辑实现无中断的监听增删。
代码实现示例
核心业务服务实现
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.Sinks; import org.springframework.stereotype.Service; import java.util.concurrent.ConcurrentHashMap; @Service public class HardwareEventService { // 动态接收新硬件监听流的Sink,多播模式+背压缓存,可根据实际场景调整选型 private final Sinks.Many<Flux<HardwareEvent>> hardwarePublishersSink = Sinks.many() .multicast() .onBackpressureBuffer(); // 维护硬件ID和对应监听终止信号的映射,用于移除监听 private final ConcurrentHashMap<HardwareId, Sinks.Empty<Void>> stopSignalMap = new ConcurrentHashMap<>(); // 对外暴露的合并后统一事件流 public Flux<HardwareEvent> getMergedEventStream() { // 等价于动态merge所有上层下发的硬件监听流 return hardwarePublishersSink.asFlux() .flatMap(flux -> flux); } // 新增硬件监听 public void addHardwareListener(HardwareId id) { if (stopSignalMap.containsKey(id)) { // 避免重复注册同个硬件的监听 return; } Sinks.Empty<Void> stopSink = Sinks.empty(); stopSignalMap.put(id, stopSink); // 构造带终止逻辑的硬件监听流 Flux<HardwareEvent> hardwareEventFlux = listenHardware(id) .takeUntilOther(stopSink.asMono()) .doFinally(signalType -> stopSignalMap.remove(id)); // 流结束后自动清理映射 // 注入新的监听流到合并容器 hardwarePublishersSink.tryEmitNext(hardwareEventFlux) .orElseThrow(); // 发射失败可自定义异常处理逻辑 } // 移除硬件监听 public void removeHardwareListener(HardwareId id) { Sinks.Empty<Void> stopSink = stopSignalMap.remove(id); if (stopSink != null) { // 发送终止信号结束对应硬件的监听流,不影响其他流和总合并流 stopSink.tryEmitEmpty() .orElseThrow(); } } // 你原有的硬件监听方法,返回无限事件流 private Flux<HardwareEvent> listenHardware(HardwareId id) { // 原有业务实现 } }
接口层实现
import org.springframework.http.MediaType; import org.springframework.web.bind.annotation.*; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @RestController public class HardwareEventController { private final HardwareEventService hardwareEventService; public HardwareEventController(HardwareEventService hardwareEventService) { this.hardwareEventService = hardwareEventService; } // 事件订阅接口,返回SSE流适配前端长连接场景 @GetMapping(value = "/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<HardwareEvent> subscribeEvents() { return hardwareEventService.getMergedEventStream(); } // 新增硬件监听接口 @PostMapping("/listeners/{hardwareId}") public Mono<Void> addListener(@PathVariable HardwareId hardwareId) { hardwareEventService.addHardwareListener(hardwareId); return Mono.empty(); } // 移除硬件监听接口 @DeleteMapping("/listeners/{hardwareId}") public Mono<Void> removeListener(@PathVariable HardwareId hardwareId) { hardwareEventService.removeHardwareListener(hardwareId); return Mono.empty(); } }
注意事项
- Sink选型可根据业务调整:如果要求新订阅者能收到所有历史注册的硬件事件流,可将
hardwarePublishersSink替换为Sinks.many().replay().all() - 背压策略可根据事件生产速率调整,可替换为
onBackpressureDrop、onBackpressureLatest等其他策略 - 若
listenHardware方法持有外部资源(如TCP连接、硬件端口),可在流的doOnCancel、doFinally钩子中添加资源释放逻辑,确保移除监听时资源正确回收 - 示例中的
orElseThrow可替换为自定义的失败处理逻辑,如打日志、重试等
内容的提问来源于stack exchange,提问作者Jiinxy
相关产品推荐
相关产品推荐

