如何取消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
相关产品推荐
相关产品推荐

