如何在JDK 9 Flow API中实现Flow.Subscriber取消订阅?
这个问题我之前在项目里碰到过,JDK 9 Flow API的取消订阅逻辑其实藏在Subscription对象里,虽然subscribe()方法本身没提供直接取消的入口,但咱们可以通过合理的状态管理来搞定,具体步骤如下:
第一步:在Room类中维护用户与订阅的映射
首先需要一个线程安全的容器来记录每个用户对应的Subscription实例,因为用户进入/离开房间可能是并发操作,推荐用ConcurrentHashMap:private final Map<User, Flow.Subscription> userSubscriptions = new ConcurrentHashMap<>();第二步:为每个用户创建专属的Subscriber代理
不要直接让Room本身作为Subscriber订阅User(否则无法区分不同User的Subscription),而是为每个进入房间的用户创建一个内部Subscriber代理,既关联用户与Subscription,又把事件转发给Room处理:public void userEnter(User user) { user.subscribe(new Flow.Subscriber<Notification>() { @Override public void onSubscribe(Flow.Subscription subscription) { // 绑定用户与对应的订阅关系 userSubscriptions.put(user, subscription); // 根据业务需求请求数据,比如请求无限量事件: subscription.request(Long.MAX_VALUE); } @Override public void onNext(Notification item) { // 把事件转发给Room的核心处理逻辑 Room.this.handleUserNotification(user, item); } @Override public void onError(Throwable throwable) { // 发生错误时自动清理订阅关系 userSubscriptions.remove(user); Room.this.handleUserError(user, throwable); } @Override public void onComplete() { // 用户发布端完成时清理订阅 userSubscriptions.remove(user); Room.this.handleUserComplete(user); } }); }这里的
handleUserNotification、handleUserError就是你原来Room类中处理事件的核心方法,可以直接复用之前的逻辑。第三步:用户离开时取消订阅
当用户离开房间时,从映射中取出对应的Subscription,调用它的cancel()方法就能终止订阅关系,最后清理映射:public void userLeave(User user) { Flow.Subscription subscription = userSubscriptions.remove(user); if (subscription != null) { // 调用cancel()终止订阅,停止接收该用户的事件 subscription.cancel(); } }
关键原理说明
JDK 9 Flow API的订阅关系是通过Subscription实例唯一标识的,每个subscribe()调用都会生成一个独立的Subscription。调用subscription.cancel()后,Publisher就会停止向该Subscriber发送事件,同时释放相关资源,完美解决用户离开后的订阅清理问题。
内容的提问来源于stack exchange,提问作者António Almeida

