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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 17:53:20