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

如何取消Vert.x中基于ReactiveStreams发布SSE事件的订阅?

如何取消Vert.x中Reactive Streams的SSE数据流订阅

当然可以取消这个Reactive Streams的订阅!在Vert.x场景下,核心是利用Reactive Streams规范中的Subscription对象来控制订阅生命周期,结合SSE的HTTP上下文,你可以实现自动取消(比如连接关闭时)或者手动触发取消。下面是具体的实现方案:

核心原理

当你订阅Reactive Streams的Publisher时,Subscriber的onSubscribe方法会收到一个Subscription实例。调用这个实例的cancel()方法,就能立即终止订阅,Publisher会停止发送数据并清理相关资源(比如你的periodicStream定时器会被自动停止)。

方案1:HTTP连接关闭时自动取消

在SSE场景中,最常见的需求是当客户端断开连接时自动停止数据流。你可以把Subscription和HTTP请求上下文绑定,在请求结束时触发取消:

// 先定义你的SSE Publisher
Publisher<String> ssePublisher = ReactiveStreams.fromPublisher(vertx.periodicStream(1000).toPublisher())
    .map(l -> String.format("Number of Customer added %s .%n", customerRepository.findAll().size()))
    .buildRs();

// 配置SSE路由
router.get("/sse").handler(ctx -> {
    HttpServerResponse response = ctx.response();
    // 设置SSE必要的响应头
    response.putHeader("Content-Type", "text/event-stream");
    response.putHeader("Cache-Control", "no-cache");
    response.putHeader("Connection", "keep-alive");
    response.setChunked(true);

    // 订阅Publisher并保存Subscription
    ssePublisher.subscribe(new Subscriber<>() {
        private Subscription subscription;

        @Override
        public void onSubscribe(Subscription s) {
            this.subscription = s;
            // 请求所有可用数据(按需调整,比如每次请求1条)
            s.request(Long.MAX_VALUE);

            // 当HTTP连接关闭时自动取消订阅
            ctx.request().endHandler(v -> {
                subscription.cancel();
                System.out.println("SSE订阅已随连接关闭取消");
            });
        }

        @Override
        public void onNext(String s) {
            // 发送SSE事件格式的数据
            response.write("data: " + s + "\n\n");
        }

        @Override
        public void onError(Throwable t) {
            // 发生错误时关闭响应
            response.end();
            t.printStackTrace();
        }

        @Override
        public void onComplete() {
            // 数据流完成时关闭响应
            response.end();
        }
    });
});

方案2:手动触发取消(比如通过API)

如果你需要主动触发取消(比如给特定客户端停止推送),可以把Subscription存储在一个线程安全的容器中,通过额外的API来调用cancel():

// 用ConcurrentHashMap存储活跃订阅,确保线程安全
private final ConcurrentHashMap<String, Subscription> activeSubscriptions = new ConcurrentHashMap<>();

// SSE路由,带客户端ID参数
router.get("/sse/:clientId").handler(ctx -> {
    String clientId = ctx.pathParam("clientId");
    HttpServerResponse response = ctx.response();
    // 设置SSE响应头...(同方案1)
    response.setChunked(true);

    ssePublisher.subscribe(new Subscriber<>() {
        private Subscription subscription;

        @Override
        public void onSubscribe(Subscription s) {
            this.subscription = s;
            // 将订阅存入Map,关联客户端ID
            activeSubscriptions.put(clientId, s);
            s.request(Long.MAX_VALUE);

            // 连接关闭时自动清理订阅
            ctx.request().endHandler(v -> {
                activeSubscriptions.remove(clientId);
                subscription.cancel();
            });
        }

        // onNext、onError、onComplete实现同方案1...
    });
});

// 手动取消订阅的API
router.post("/cancel-sse/:clientId").handler(ctx -> {
    String clientId = ctx.pathParam("clientId");
    Subscription subscription = activeSubscriptions.remove(clientId);
    
    if (subscription != null) {
        subscription.cancel();
        ctx.response().setStatusCode(200).end("客户端[" + clientId + "]的SSE订阅已取消");
    } else {
        ctx.response().setStatusCode(404).end("未找到客户端[" + clientId + "]的活跃订阅");
    }
});

注意事项

  • 资源清理:调用subscription.cancel()后,Vert.x的periodicStream会自动停止定时器,不会继续占用资源。
  • 线程安全:存储多个订阅时,一定要用线程安全的容器(比如ConcurrentHashMap),避免并发访问问题。
  • 请求策略:request(Long.MAX_VALUE)是一次性请求所有数据,如果你需要更精细的流量控制,可以按需调整请求的数量(比如每次请求1条,处理完再请求下一条)。

内容的提问来源于stack exchange,提问作者Francesco Marchioni

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:13:55