Flow / SubmissionPublisher / Reactive Streams 源码
概述
JDK 9 引入 java.util.concurrent.Flow,把 Reactive Streams(响应式流) 标准内建进 JDK:Publisher 发布、Subscriber 订阅、Subscription 控制背压、Processor 变换。SubmissionPublisher 则是该标准的参考实现——一个开箱即用的异步发布器,把 submit() 的元素异步推送给订阅者,并通过需求计数(demand)实现背压。
本文基于 OpenJDK 21 源码,从 Flow 的接口契约出发,拆解 SubmissionPublisher 的发布分发、BufferedSubscription 的队列与需求控制、异步执行与消费流程。
核心源码解析
① Flow 的 4 个接口
java
public final class Flow {
private Flow() {} // 纯工具容器,禁止实例化
@FunctionalInterface
public static interface Publisher<T> {
public void subscribe(Subscriber<? super T> subscriber);
}
public static interface Subscriber<T> {
public void onSubscribe(Subscription subscription); // 订阅建立回调
public void onNext(T item); // 数据推送
public void onError(Throwable throwable); // 错误通知
public void onComplete(); // 完成通知
}
public static interface Subscription {
public void request(long n); // 请求 n 个元素(背压)
public void cancel(); // 取消订阅
}
public static interface Processor<T, R> extends Subscriber<T>, Publisher<R> {
// 同时是订阅者与发布者,用于转换
}
}- 完整的 Reactive Streams 标准:四个接口恰好对应规范中的四类参与者;
Subscriber的四个回调构成生命周期(onSubscribe→onNext*→onComplete/onError)。 - 背压协议:
Subscription.request(n)是唯一合法的"请继续推"信号——发布者最多推送n个元素,超出的必须缓存或丢弃,订阅者不会被淹没。 - 契约约束:
onNext必须按请求串行调用、request可累加、cancel后不得再推送、所有回调不得阻塞。SubmissionPublisher保证了这些约束的实现。
② SubmissionPublisher.submit(T item) 的发布流程
java
public class SubmissionPublisher<T> implements Publisher<T>, AutoCloseable {
final Executor executor; // 异步执行器
final BiConsumer<Subscriber<? super T>, Throwable> onNextHandler; // 异常处理器
final Set<BufferedSubscription<T>> subscribers = new HashSet<>();
int maxBufferCapacity; // 最大缓冲容量(默认 256)
volatile boolean closed;
volatile Throwable closedException;
public int submit(T item) {
return doOffer(item, Long.MAX_VALUE, null); // 无超时等待
}
public int offer(T item, BiPredicate<Subscriber<? super T>, ? super T> onDrop) {
return doOffer(item, 0L, onDrop); // 立即可用
}
private int doOffer(T item, long nanos, BiPredicate<...> onDrop) {
// ① 遍历订阅者集合
for (BufferedSubscription<T> s : subscribers) {
// ② 尝试入队,超时/失败时按 onDrop 处理
if (!s.offer(item, nanos)) {
if (onDrop == null || !onDrop.test(s.subscriber, item)) {
s.cancel(); // 满则取消该订阅
}
}
}
return subscribers.size(); // 返回尝试分发的订阅者数
}
}submit把元素异步推给所有订阅者:内部调用doOffer,每个订阅者通过其BufferedSubscription.offer入队并唤醒消费线程。- 满缓冲策略:队列满时
offer返回 false,若未提供onDrop谓词则取消该订阅者;offer(item, onDrop)给调用方机会决定丢弃或降级。 - 返回值是订阅者数量(被推送成功的目标数),方便调用方统计;
submit的返回语义为"至少尝试了 N 个订阅者"。
③ SubmissionPublisher 的 BufferedSubscription 内部队列
java
static final class BufferedSubscription<T> extends SubmissionPublisher<T>
implements Subscription, Runnable {
final Subscriber<? super T> subscriber;
final Executor executor;
final BiConsumer<Subscriber<? super T>, Throwable> onNextHandler;
final AtomicLong demand = new AtomicLong(); // 需求计数
final ConcurrentLinkedQueue<T> queue = new ConcurrentLinkedQueue<>(); // 内部队列
volatile boolean active; // 消费任务是否在跑
...
}ConcurrentLinkedQueue+AtomicLong:队列无锁存储元素,demand记录订阅者尚未消费的请求量;两者组合实现背压——onNext只能在demand > 0时触发。- 继承自
SubmissionPublisher(JDK 实现技巧):BufferedSubscription复用发布器字段但以"单订阅者"形态工作,减少代码重复。 - 消费由
Runnable触发:executor.execute(this)把消费任务提交到执行器,active标志保证同一时刻只有一个消费循环在跑(避免并发重复消费)。
④ Subscription.request(long n) 的背压
java
public void request(long n) {
if (n <= 0) throw new IllegalArgumentException("non-positive subscription request");
demand.getAndAdd(n); // ① 累加需求计数
drainLoop(); // ② 尝试立即消费
}
private void drainLoop() {
// 已在消费则返回,保证单线程消费
if (active || !casActive()) return;
int missed = 1;
for (;;) {
long demand = this.demand.get();
long c = 0;
while (c < demand) {
T item = queue.poll();
if (item == null) break; // 队列空,等待下次
try {
subscriber.onNext(item); // 消费一个:回调订阅者
} catch (Throwable ex) {
handleOnNextError(ex, item); // 异常走 onError 通道
return;
}
c++;
}
if (c == demand && (demand = this.demand.get()) == 0) break; // 无新请求
...
}
}addAndGet(n):request累加需求,支持多次调用叠加(request(1)多次等价于一次request(3))。drainLoop消费:从队列拉取元素并回调onNext,数量不超过当前demand;消费完一轮后检查是否有新请求/新元素,用missed计数处理竞争。- 背压闭环:
request(n)决定推送上限 → 队列缓存超出部分 →onNext只按需求触发 → 订阅者通过控制request频率调节速率。这是响应式流不会压垮消费者的核心。
⑤ SubmissionPublisher 的 Executor 选择
java
public SubmissionPublisher() {
this(AsyncExecution.getExecutor(), BUFFER, null); // BUFFER = 32(每订阅者缓冲)
}
// jdk.internal.util.concurrent.AsyncExecution
static Executor getExecutor() {
return new DelegateExecutor(ThreadLocalRandom.current().nextBoolean()
? new ForkJoinPool() // 自建 ForkJoinPool
: ForkJoinPool.commonPool()); // 公共线程池(默认)
}- 默认
ForkJoinPool.commonPool():与CompletableFuture共享的公共池,异步消费onNext()在其中执行,submit()立即返回不阻塞调用线程。 - 为什么异步:
onNext在drainLoop中于执行器线程回调,发布者线程只做入队——发布与消费彻底解耦,即使订阅者处理慢也不会拖慢submit。 - 可通过构造器
SubmissionPublisher(Executor, int, BiConsumer)自定义执行器(如单线程池保证顺序、或限流池);maxBufferCapacity决定每个订阅者的队列上限。 - 注意:
commonPool的并行度受-Djava.util.concurrent.ForkJoinPool.common.parallelism影响,高吞吐场景可显式指定执行器。
⑥ SubmissionPublisher.consume(Consumer) 的便捷消费
java
public void consume(Consumer<? super T> consumer) {
subscribe(new ConsumerSubscriber<T>(consumer));
}ConsumerSubscriber是Subscriber的适配实现:收到onSubscribe立即request(Long.MAX_VALUE)(无界请求),随后每个onNext自动调用consumer.accept(item)。- 对调用方而言,
consume(c)一行即可订阅并消费,省去手动管理Subscription的样板代码。 - 局限性:无界请求意味着不做背压(
Long.MAX_VALUE),只适合消费速度远快于生产速度的场景;需要节流时应使用标准subscribe+ 手动request。 consume是"单次使用"便捷入口(内部consumed标志防重复),结合close()可编排"生产-消费"生命周期。
⑦ Flow.Processor 的转换器
java
public static interface Processor<T, R> extends Subscriber<T>, Publisher<R> {
}- 双重身份:
Processor同时实现Subscriber<T>(上游向它推送)与Publisher<R>(它向下游发布),是响应式链路的中间变换节点。 - 责任链组装:典型用法是
Processor包装转换逻辑,链式连接上游与下游——publisher.subscribe(processor)后processor.subscribe(downstream),数据流经"上游 → 处理器 → 下游"逐级传递。 SubmissionPublisher本身可作为处理器骨架:继承它并重写onNext(transform 后submit(result)),即得到带背压的转换节点;java.util.concurrent未提供现成Processor实现,通常由第三方库(Reactor、RxJava)提供。
总结
| 组件 | 职责 | 核心机制 |
|---|---|---|
Flow 4 接口 | Reactive Streams 契约 | 订阅/推送/背压/变换四角色 |
submit / offer | 元素分发 | doOffer 遍历订阅者入队 |
BufferedSubscription | 单订阅者缓冲 | ConcurrentLinkedQueue + AtomicLong demand |
request(n) | 背压调节 | addAndGet 累加 + drainLoop 限量消费 |
executor | 异步执行 | 默认 ForkJoinPool.commonPool |
consume | 便捷消费 | ConsumerSubscriber 无界请求 |
Processor | 链路变换 | Subscriber + Publisher 双重身份 |
SubmissionPublisher 的价值在于把背压协议落地成一套可用的线程模型:入队异步、消费限速、失败隔离。理解 demand 计数与 drainLoop 单线程消费循环,就抓住了响应式背压的实现本质。