如何正确取消WebFlux创建的SSE流?
我来分享一个完整的解决方案,把遗留系统的监听器回调转换成WebFlux的无限数据流,包含服务端和客户端的可运行代码:
服务端实现
首先我们模拟遗留系统的SampleEmitter,它支持注册监听器,会定时生成数据并通知监听器:
import java.util.ArrayList; import java.util.List; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.function.Consumer; public class SampleEmitter { private final List<Consumer<String>> listeners = new ArrayList<>(); private final ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); private int counter = 0; public SampleEmitter() { // 每1秒生成一条数据 executor.scheduleAtFixedRate(() -> { String data = "Generated data: " + counter++; listeners.forEach(listener -> listener.accept(data)); }, 0, 1, TimeUnit.SECONDS); } public void registerListener(Consumer<String> listener) { listeners.add(listener); } public void removeListener(Consumer<String> listener) { listeners.remove(listener); } public void shutdown() { executor.shutdown(); } }
接下来是核心部分:把监听器的回调转换成WebFlux的Flux。我们用Flux.create来创建可推送的数据流,并且在订阅取消时移除监听器,避免内存泄漏:
import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.Flux; @RestController public class DataStreamController { private final SampleEmitter sampleEmitter = new SampleEmitter(); @GetMapping("/stream") public Flux<String> streamData() { return Flux.create(sink -> { // 创建监听器,收到数据就推送到Flux Consumer<String> listener = data -> { if (!sink.isCancelled()) { sink.next(data); } }; // 注册监听器到遗留系统 sampleEmitter.registerListener(listener); // 当订阅取消时,移除监听器,清理资源 sink.onCancel(() -> sampleEmitter.removeListener(listener)); sink.onDispose(sampleEmitter::shutdown); }); } }
客户端实现
客户端用WebClient来订阅这个无限数据流,处理每个数据,并且可以在需要时取消订阅:
import org.springframework.web.reactive.function.client.WebClient; import reactor.core.Disposable; import java.util.concurrent.TimeUnit; public class StreamClient { public static void main(String[] args) throws InterruptedException { WebClient webClient = WebClient.create("http://localhost:8080"); // 订阅数据流 Disposable subscription = webClient.get() .uri("/stream") .retrieve() .bodyToFlux(String.class) .doOnNext(data -> System.out.println("Received: " + data)) .doOnCancel(() -> System.out.println("Subscription cancelled")) .subscribe(); // 模拟订阅5秒后取消 TimeUnit.SECONDS.sleep(5); subscription.dispose(); } }
关键注意点
- Flux.create的适配:这是将传统监听器模式转换为响应式流的最优方案,它支持主动推送数据,同时能监听订阅的生命周期变化。
- 取消订阅的资源清理:务必在
sink.onCancel()中移除遗留系统的监听器,否则即使客户端取消订阅,遗留系统仍会继续推送数据,引发不必要的资源消耗。 - 无限流的持续特性:只要服务端处于运行状态,数据流就会不断生成新数据,客户端可以灵活控制订阅的启停。
内容的提问来源于stack exchange,提问作者countryroadscat
相关产品推荐
相关产品推荐

