You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 05:42:03