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

如何正确取消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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 11:59:03