基于Flow API实现观察者模式:代码完善与合理性咨询
问题解答
1. 完善onSubscribe与subscribe方法
Java Flow API基于请求-响应模型,必须正确实现订阅逻辑和流量控制,以下是修正后的核心代码:
修正ObservableDouble
import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.Flow; import java.util.concurrent.atomic.AtomicReference; public class ObservableDouble implements Flow.Publisher<Double> { private final AtomicReference<Double> value = new AtomicReference<>(Double.MIN_VALUE); // 使用CopyOnWriteArrayList保证多线程下订阅者列表的安全 private final List<Flow.Subscriber<? super Double>> subscribers = new CopyOnWriteArrayList<>(); public double getValue() { return value.get(); } public void setValue(double newValue) { Double oldValue = value.getAndSet(newValue); // 仅当值真正变化时才推送,避免无效更新 if (!oldValue.equals(newValue)) { subscribers.forEach(subscriber -> { subscriber.onNext(newValue); }); } } public boolean isSet() { return value.get() != Double.MIN_VALUE; } @Override public void subscribe(Flow.Subscriber<? super Double> subscriber) { if (subscriber == null) { throw new NullPointerException("Subscriber cannot be null"); } // 创建自定义Subscription,处理请求和取消逻辑 SubscriptionImpl subscription = new SubscriptionImpl(subscriber); subscriber.onSubscribe(subscription); subscribers.add(subscriber); } // 内部类实现Subscription,处理流量控制 private class SubscriptionImpl implements Flow.Subscription { private final Flow.Subscriber<? super Double> subscriber; private volatile boolean cancelled; private final java.util.concurrent.atomic.AtomicLong requested = new java.util.concurrent.atomic.AtomicLong(0); public SubscriptionImpl(Flow.Subscriber<? super Double> subscriber) { this.subscriber = subscriber; } @Override public void request(long n) { if (n <= 0) { subscriber.onError(new IllegalArgumentException("Request must be positive")); return; } requested.addAndGet(n); // 若当前有已设置的值,立即推送当前值 if (isSet() && !cancelled) { subscriber.onNext(value.get()); } } @Override public void cancel() { cancelled = true; subscribers.remove(subscriber); } } }
修正ObserverDouble
import java.util.concurrent.Flow; import java.util.function.Consumer; public class ObserverDouble implements Flow.Subscriber<Double> { private Consumer<Double> consumer; private Flow.Subscription subscription; public void setConsumer(Consumer<Double> consumer) { this.consumer = consumer; } @Override public void onSubscribe(Flow.Subscription subscription) { this.subscription = subscription; // 请求无限量数据,也可根据需求设置具体数值实现流量控制 subscription.request(Long.MAX_VALUE); } @Override public void onNext(Double item) { if (consumer != null) { consumer.accept(item); } } @Override public void onError(Throwable throwable) { throwable.printStackTrace(); } @Override public void onComplete() { // 可添加完成后的清理逻辑 } }
2. Consumer的合理性分析
使用Consumer<Double>是合理且推荐的设计,原因如下:
- 简化回调逻辑:无需让用户每次都实现完整的
Flow.Subscriber接口,只需传入一个简单的消费函数即可处理值更新。 - 灵活性高:可以随时更换消费逻辑,无需创建新的
Subscriber实例。 - 符合函数式编程风格:Java 8+的函数式接口能让代码更简洁易读。
需要注意的细节:
- 必须在
onNext中做非空检查,避免consumer为null时抛出NullPointerException。 - 若需要处理错误或完成事件,可以扩展
ObserverDouble,添加Consumer<Throwable>和Runnable分别处理onError和onComplete逻辑,进一步增强灵活性。
内容的提问来源于stack exchange,提问作者Paul
相关产品推荐
相关产品推荐

