Java 基础体系 · 第 70/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java Flow 与响应式流:Publisher、Subscriber、背压和协议正确性

java.util.concurrent.Flow 是 Java 标准库对响应式流协议(Reactive Streams)的接口定义。它解决的不是“如何异步执行一个任务”,而是如何让一个数据生产者和一个数据消费者在异步、跨线程、不同处理速度的情况下,仍然按照明确的协议传递数据。

这个协议的核心约束是:

  1. Subscriber 必须先收到 onSubscribe
  2. Subscriber 通过 Subscription.request(n) 声明自己还能接收多少个元素;
  3. Publisher 发送的 onNext 数量不能超过已请求数量;
  4. onErroronComplete 是终止信号,终止后不能再发送其他信号;
  5. 同一个订阅上的信号必须保持顺序,不能因为并发发送而破坏协议。

因此,响应式流的难点不在接口数量,而在于需求量、并发、取消、缓存和终止状态必须同时保持一致

Flow API 与响应式流的关系

Java 9 引入了 java.util.concurrent.Flow,Java 25 LTS 中仍然使用这组 API。它包含四个核心接口:

public final class Flow {
    public interface Publisher<T> {
        void subscribe(Subscriber<? super T> subscriber);
    }

    public interface Subscriber<T> {
        void onSubscribe(Subscription subscription);
        void onNext(T item);
        void onError(Throwable throwable);
        void onComplete();
    }

    public interface Subscription {
        void request(long n);
        void cancel();
    }

    public interface Processor<T, R>
            extends Subscriber<T>, Publisher<R> {
    }
}

它们分别表示:

  • Publisher<T>:可以发布 T 类型数据的组件;
  • Subscriber<T>:接收 T 类型数据的组件;
  • Subscription:一次具体订阅关系的控制句柄;
  • Processor<T, R>:同时是 Subscriber<T>Publisher<R>,可以接收 T 并发布 R

Flow 的接口和 Reactive Streams 规范采用相同的核心协议模型。Reactive Streams 还定义了 org.reactivestreams.Publisher 等接口,它们不属于 JDK。第三方响应式库通常提供两者之间的适配器,但 Java 标准库本身并不自动提供所有外部库的适配实现。

Flow 不是一个完整的响应式编程框架。它没有内置 mapfilter、线程调度器、重试、合并等高层操作符,只定义了数据流参与者之间必须遵守的底层契约。

一次订阅包含哪些角色

调用:

publisher.subscribe(subscriber);

并不只是注册一个回调。它会建立一条独立的订阅通道,典型流程如下:

sequenceDiagram
    participant P as Publisher
    participant S as Subscriber
    participant Q as Subscription

    S->>P: subscribe(S)
    P->>S: onSubscribe(Q)
    S->>Q: request(2)
    P->>S: onNext(item-1)
    P->>S: onNext(item-2)
    S->>Q: request(1)
    P->>S: onNext(item-3)
    S->>Q: cancel()
    P-->>S: 停止继续发送

这里有一个重要区别:

  • Publisher 表示“可以发布数据的能力”;
  • Subscriber 表示“如何处理数据”;
  • Subscription 表示“这一次连接的状态”,包括需求量和取消状态。

同一个 Publisher 可以被多个 Subscriber 订阅。每次订阅通常都应有独立的 Subscription,不能把所有订阅者共享成一个全局需求计数器,否则一个订阅者的 request 会错误地影响其他订阅者。

Subscriber 的生命周期

一个合法的订阅生命周期必须从 onSubscribe 开始:

onSubscribe
    ├── onNext
    ├── onNext
    ├── ...
    └── onComplete

或者:

onSubscribe
    ├── onNext
    ├── ...
    └── onError

onCompleteonError 都是终止信号。一条订阅只能有一个终止结果:

  • 可以正常完成;
  • 可以异常结束;
  • 可以被取消而不再收到终止回调;
  • 不能先 onCompleteonError
  • 不能先 onErroronNext

cancel()onComplete() 不是同一个概念:

  • onComplete() 表示 Publisher 已经正常结束;
  • cancel() 表示 Subscriber 不再希望继续接收;
  • 取消后,通常不会要求 Publisher 再向该订阅发送 onComplete
  • 取消是异步协议中的停止请求,已经在其他线程进入调用过程的信号可能存在竞态,因此实现必须正确处理并发停止。

Subscriber 应当把 onSubscribe 保存下来,并通过它进行需求控制和取消:

final class LoggingSubscriber<T> implements Flow.Subscriber<T> {
    private Flow.Subscription subscription;

    @Override
    public void onSubscribe(Flow.Subscription subscription) {
        this.subscription = subscription;
        subscription.request(1);
    }

    @Override
    public void onNext(T item) {
        System.out.println("item = " + item);
        subscription.request(1);
    }

    @Override
    public void onError(Throwable throwable) {
        throwable.printStackTrace();
    }

    @Override
    public void onComplete() {
        System.out.println("complete");
    }
}

这段代码采用“一次请求一个”的策略。收到一个元素、处理完成后,再请求下一个元素。它简单地体现了背压,但不一定有最高吞吐量,因为每个元素都需要一次需求更新。

背压:需求量如何约束发送量

背压(backpressure)是消费者反向约束生产者发送速度的机制。

设某个 Subscription 的需求量为 D

  • request(n) 中的 n 必须是正数;
  • 调用 request(n) 后,需求量增加 n
  • 每发送一个 onNext,需求量减少 1;
  • Publisher 发送的元素数不能超过需求量;
  • request(Long.MAX_VALUE) 通常表示近似无限需求。

可以用如下不变量表示:

已发送但尚未被需求覆盖的元素数 ≤ 已请求但尚未消耗的需求量

更形式化地,设:

  • R:累计请求量;
  • E:累计发送的 onNext 数量;
  • C:是否已经取消或终止。

在没有溢出处理的简化模型中,合法状态必须满足:

E ≤ R

每次事件的状态变化为:

request(n), n > 0:
    R := R + n

onNext(item):
    发送前必须满足 E < R
    发送后 E := E + 1

cancel():
    C := true
    后续不得继续正常发送

onComplete/onError:
    进入终止状态

例如,初始状态为:

R = 0, E = 0

调用 request(2)

R = 2, E = 0

发送第一个元素:

R = 2, E = 1

发送第二个元素:

R = 2, E = 2

此时 E == R,继续发送第三个元素就是协议错误。Subscriber 如果还想接收数据,必须再次调用:

subscription.request(1);

随后:

R = 3, E = 2

于是又允许发送一个元素。

需求量实际表示的是“允许继续发送的额度”,不是已经处理完成的元素数。一个 Subscriber 可以提前请求 100 个元素,然后在内部队列中逐个处理。

request(n) 不是“立即发送 n 个元素”

request(n) 是一个需求信号,不是同步调用命令。它不保证:

  • Publisher 立即发送;
  • Publisher 一定能发送恰好 n 个;
  • onNext 会在调用 request 的线程执行;
  • 所有请求都能最终满足。

例如,Publisher 可能只有三个元素,而 Subscriber 请求了 100 个。它可以发送三个元素,然后调用 onComplete。需求量没有被“填满”的义务,完成信号可以提前结束数据流。

同样,Publisher 也可以暂时没有数据。请求量可以保持为正数,等数据将来产生时再发送。

请求量的溢出处理

如果实现直接执行:

requested += n;

就可能发生 long 溢出。例如:

requested = Long.MAX_VALUE - 1
n         = 10

结果会变成负数,进而破坏需求量判断。

常见的规范实现会把需求量饱和到 Long.MAX_VALUE

static long addCap(long current, long n) {
    long next = current + n;
    return next < 0L ? Long.MAX_VALUE : next;
}

这段代码只适用于 currentn 都已经确认是正数的场景。生产级实现还需要在原子更新或锁保护下执行,否则多个线程同时调用 request 仍可能丢失更新。

Long.MAX_VALUE 通常被当作无界需求。例如:

subscription.request(Long.MAX_VALUE);

这并不意味着内存可以无限增长,也不代表 Publisher 必须永远运行。它只是表示 Subscriber 不再要求 Publisher 按有限额度限制发送。数据源本身仍可能耗尽、失败或被取消。

一个可运行的 Flow 示例

下面的程序使用 JDK 自带的 SubmissionPublisher,展示:

  • Publisher 如何发布数据;
  • Subscriber 如何在 onSubscribe 中请求数据;
  • onNext 处理完成后如何继续请求;
  • Publisher 如何正常关闭;
  • Subscriber 如何接收 onComplete
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Flow;
import java.util.concurrent.SubmissionPublisher;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class FlowDemo {
    public static void main(String[] args) throws InterruptedException {
        ExecutorService executor =
                Executors.newVirtualThreadPerTaskExecutor();

        CountDownLatch done = new CountDownLatch(1);

        try (SubmissionPublisher<Integer> publisher =
                     new SubmissionPublisher<>(executor, 2)) {

            publisher.subscribe(new Flow.Subscriber<>() {
                private Flow.Subscription subscription;

                @Override
                public void onSubscribe(Flow.Subscription subscription) {
                    this.subscription = subscription;

                    // 先只允许发送一个元素
                    subscription.request(1);
                }

                @Override
                public void onNext(Integer item) {
                    System.out.println(
                            Thread.currentThread().getName()
                                    + " received " + item);

                    // 模拟消费耗时
                    try {
                        Thread.sleep(50);
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                        subscription.cancel();
                        return;
                    }

                    // 当前元素处理完成后,再允许发送一个
                    subscription.request(1);
                }

                @Override
                public void onError(Throwable throwable) {
                    throwable.printStackTrace();
                    done.countDown();
                }

                @Override
                public void onComplete() {
                    System.out.println("complete");
                    done.countDown();
                }
            });

            for (int i = 1; i <= 5; i++) {
                publisher.submit(i);
                System.out.println("submitted " + i);
            }

            // 关闭表示不再提交新元素。
            // 已经提交的元素仍会按协议发送完,然后触发 onComplete。
            publisher.close();

            done.await();
        } finally {
            executor.close();
        }
    }
}

可以使用 Java 25 编译和运行:

javac FlowDemo.java
java FlowDemo

输出顺序会受到异步调度影响,但数据顺序应保持为 15,并最终出现:

complete

submittedreceived 的相对顺序不固定,因为 SubmissionPublisher 使用异步执行器投递信号。submitted 5 可能先于 received 1 打印。

这里有三个容易混淆的地方。

第一,submit 不等于 onNext

publisher.submit(i) 表示把元素提交给 SubmissionPublisher。它可能先进入内部缓冲区,之后才由异步任务调用 Subscriber 的 onNext

因此:

submit(1)

不表示:

onNext(1) 已经执行完毕

第二,背压可能表现为提交方阻塞

示例把每个订阅者的缓冲容量设置为 2:

new SubmissionPublisher<>(executor, 2)

Subscriber 每次只请求一个元素,而且每个元素处理 50 毫秒。如果生产者提交速度更快,缓冲区可能被填满。SubmissionPublisher.submit 在这种情况下可能阻塞提交线程,等待订阅者继续消费。

这说明背压并不总是表现为“生产者收到一个显式的暂停回调”。一种实现可以通过阻塞、排队、丢弃、失败或降低上游生成速度来体现背压,具体取决于组件设计。

第三,关闭 Publisher 不是取消 Subscriber

publisher.close();

表示 Publisher 不再接受新的提交,并在已提交数据处理完成后正常结束订阅。

如果 Subscriber 不想继续接收,应调用:

subscription.cancel();

取消是订阅者主动终止这条订阅通道,而不是让整个 Publisher 对所有订阅者关闭。

没有请求时会发生什么

下面的 Subscriber 在 onSubscribe 中不调用 request

final class NeverRequestSubscriber<T>
        implements Flow.Subscriber<T> {

    @Override
    public void onSubscribe(Flow.Subscription subscription) {
        System.out.println("subscribed");
        // 故意不 request
    }

    @Override
    public void onNext(T item) {
        System.out.println("item = " + item);
    }

    @Override
    public void onError(Throwable throwable) {
        throwable.printStackTrace();
    }

    @Override
    public void onComplete() {
        System.out.println("complete");
    }
}

合法的 Publisher 不应调用它的 onNext。因为需求量仍为 0,发送任何元素都会违反背压约束。

对于 SubmissionPublisher,元素可能暂时留在内部缓冲区,提交方最终因为缓冲区满而受到阻塞。Publisher 不会因为 Subscriber 没有请求就自动获得发送许可。

这也是一个常见误解:

“异步发送”不等于“可以无限发送”。

异步只说明生产和消费可以在不同线程进行;背压仍然限制着合法发送量。

request(0) 和负数请求

以下代码是错误的:

subscription.request(0);

以及:

subscription.request(-1);

需求量必须是正数。标准实现通常会把非法请求转化为订阅错误,例如向 Subscriber 发送:

onError(new IllegalArgumentException(...))

随后不再继续正常发送。

Subscriber 不应把 request(0) 当作“暂时暂停”的操作。Flow 没有定义这个语义。要暂时停止接收,通常有两种设计:

  1. 不再追加请求,但保留当前订阅;
  2. 调用 cancel() 永久终止订阅。

第一种只适用于 Publisher 能承受未消费需求和缓冲的情况;第二种不可逆。

回调的线程与串行性

Flow 允许异步和并发执行,但单个订阅上的信号不能随意并发交错。

错误实现可能这样发送:

executor.submit(() -> subscriber.onNext(1));
executor.submit(() -> subscriber.onNext(2));

如果两个任务同时运行,Subscriber 可能观察到:

  • onNext(2) 先于 onNext(1)
  • 两个 onNext 同时进入;
  • 一个线程执行 onComplete,另一个线程仍执行 onNext

这会破坏响应式流协议。

合法实现通常需要使用以下机制之一:

  • 单线程事件循环;
  • 每个订阅一个串行队列;
  • 锁;
  • 原子状态加无锁 drain loop;
  • 等价的串行化调度机制。

需要区分两个概念:

  • 信号串行性:同一订阅的回调不能并发或乱序;
  • 回调线程固定:回调是否始终在同一个线程执行。

协议通常要求前者,不必要求后者。一个 Publisher 可以在不同线程执行连续的 onNext,只要这些调用在逻辑上严格串行且有序。

另一方面,Subscriber 的实现仍然必须考虑线程安全。虽然同一订阅的信号应当串行,但 requestcancel 可能被业务线程、超时线程或回调线程同时调用。保存的状态、计数器和取消标志不能依赖未经保护的普通变量。

onNext 中再次 request 的问题

同步 Publisher 可能在 request(n) 调用栈内直接触发 onNext

Subscriber.onSubscribe
    -> subscription.request(1)
        -> Publisher.onNext(item)
            -> Subscriber.onNext(item)
                -> subscription.request(1)
                    -> Publisher.onNext(next)

如果数据很多,这种实现会产生深度递归,甚至导致栈溢出。

因此,Publisher 不能简单地把 request 实现成:

requested += n;
while (requested > 0 && hasNext()) {
    subscriber.onNext(next());
    requested--;
}

然后允许 onNext 内部再次进入同一个循环。更可靠的做法是使用“工作进行中”标志或队列,把重入请求转换为后续 drain:

request(n)
    增加需求量
    如果当前没有 drain:
        标记 drain 进行中
        循环发送

onNext 中再次 request(n)
    只增加需求量
    发现 drain 已经进行中
    返回,由外层 drain 继续处理

这类 drain loop 同时解决两个问题:

  1. 避免回调重入造成递归;
  2. 让需求量更新和发送过程能够被串行化。

一个需求受控的 Publisher 状态模型

对于单个订阅,可以把状态抽象为:

ACTIVE
  ├── request(n > 0) -> ACTIVE,增加 demand
  ├── onNext        -> ACTIVE,消耗 demand
  ├── cancel        -> CANCELLED
  ├── onComplete    -> COMPLETED
  └── onError       -> FAILED

CANCELLEDCOMPLETEDFAILED 都是终止状态:

终止状态
  └── 后续 request/onNext/终止信号不再产生正常数据行为

但是,状态转换必须处理并发竞争。例如:

  • 线程 A 正准备发送 onNext
  • 线程 B 同时调用 cancel()
  • 线程 A 已经通过检查,线程 B 才完成取消。

如果实现只用普通布尔值检查:

if (!cancelled) {
    subscriber.onNext(item);
}

那么检查和发送之间仍有竞态。cancel() 不是一个天然的全局中断点。实现需要通过锁、原子状态或串行队列定义明确的竞态结果。

实际协议通常允许“已经开始进行的信号”完成,但取消后不应继续无界地发送后续信号。不能把取消理解为能够撤销已经进入用户代码的 onNext

错误信号与异常传播

数据流中的失败应通过:

subscriber.onError(exception);

传播,而不是让 Publisher 在线程中静默抛出异常。

例如,一个计算型 Publisher 处理元素时失败:

try {
    R result = transform(item);
    subscriber.onNext(result);
} catch (Throwable ex) {
    subscriber.onError(ex);
}

不过,生产级实现不能只写这个 try-catch。它还必须保证:

  • onError 只发送一次;
  • onError 发送后不再 onNext
  • 已取消的订阅不再继续发送;
  • 其他订阅者是否受影响,要由 Publisher 的语义决定。

一个订阅失败,不一定要求整个多订阅者 Publisher 关闭。通常每条订阅拥有独立状态;但如果底层数据源本身失败,则所有订阅者可能都会收到错误。

Subscriber 的回调也可能抛出异常。例如:

@Override
public void onNext(Integer item) {
    throw new RuntimeException("consumer bug");
}

Subscriber 不应依赖“通过抛出异常通知 Publisher”这一方式传递业务失败。业务异常应该由 Subscriber 自己决定如何处理,或者使用更高层框架定义的错误通道。

Publisher 如果希望健壮地处理 Subscriber 回调异常,通常需要捕获异常、取消该订阅并停止后续信号。但具体清理行为属于实现责任,不能把 Subscriber 抛出的异常当作正常的 onError 业务信号。

Publisher、Subscriber 和 Processor 的数据流方向

Processor<T, R> 同时实现两个方向:

上游 Publisher<T>
        |
        v
Processor<T, R>
        |
        v
下游 Subscriber<R>

Processor 必须分别处理两套需求关系:

  1. 下游请求了多少个 R
  2. Processor 向上游请求了多少个 T

例如,一个一进一出的映射 Processor:

下游 request(5)
    -> Processor 最多向上游 request(5)
    -> 上游最多发送 5 个 T
    -> Processor 产生最多 5 个 R
    -> Processor 发送给下游

如果 Processor 一收到上游数据就继续向上游请求,而不考虑下游需求,就会形成失控缓冲:

上游高速生产
    -> Processor 无限接收
    -> Processor 内部队列增长
    -> 内存压力或延迟增加

因此,Processor 的核心不是简单地在 onNext 中转换对象,而是维护上下游需求的对应关系。对于一进一出的同步变换,常见策略是:

下游需求 = d
Processor 向上游请求 d
收到一个 T
转换为一个 R
发送一个 R
需求各减少一个

但对于过滤器、批处理器和展开器,需求关系并不相等:

  • filter 可能收到 10 个 T,只产生 2 个 R
  • buffer 可能收到 10 个 T,产生 1 个 List<T>
  • flatMap 可能收到 1 个 T,产生多个 R
  • 错误重试可能让一个上游元素被多次请求。

所以 Processor 必须根据操作符的基数变化重新设计需求传播,不能盲目把下游 request(n) 原样转发给上游。

背压不能消除所有内存风险

背压约束的是 Publisher 向 Subscriber 发送 onNext 的数量,但它不自动保证整个系统没有缓存。

一个组件可能合法地:

  1. 接收上游数据;
  2. 把数据放入自己的队列;
  3. 暂时不向下游发送;
  4. 等待下游需求增加。

这仍然可能导致队列无限增长。

因此,需要区分:

协议正确性:
    没有超过下游 request 的发送额度

资源安全性:
    队列、线程、连接、文件句柄等不会无限增长

一个无限缓冲的实现可以在协议上完全正确,却最终因为内存耗尽而失败。反过来,一个有界队列也不自动保证协议正确,因为它仍可能在没有需求时错误发送 onNext

当上游速度长期大于下游处理速度时,系统必须选择一种策略:

  • 让生产者阻塞;
  • 降低生产速度;
  • 将数据持久化到磁盘或外部队列;
  • 丢弃最新数据;
  • 丢弃最旧数据;
  • 合并或采样数据;
  • 让订阅失败。

这些是系统语义选择,不是 Flow 接口单独规定的统一答案。

SubmissionPublisher 的定位和边界

SubmissionPublisher<T> 是 JDK 提供的一个通用 Publisher 实现,适合把生产者提交的数据异步分发给多个 Subscriber。

它提供了:

SubmissionPublisher<T> publisher =
        new SubmissionPublisher<>(executor, maxBufferCapacity);

其中:

  • executor 决定异步投递任务使用的执行器;
  • maxBufferCapacity 是每个订阅者缓冲的容量上限;
  • submit 提交数据;
  • close 正常结束;
  • closeExceptionally 以异常结束。

例如:

publisher.closeExceptionally(
        new IllegalStateException("source failed"));

Subscriber 随后应收到 onError,而不是 onComplete

SubmissionPublisher 适合简单的异步发布场景,但它不是完整的响应式处理框架,也不自动提供:

  • 复杂操作符组合;
  • 数据持久化;
  • 跨进程传输;
  • 可靠消息确认;
  • 任意情况下的零缓存传输;
  • 业务级重试和死信处理。

此外,多个 Subscriber 之间的处理速度可能不同。快速 Subscriber 和慢速 Subscriber 通常有各自的缓冲与需求状态,因此慢订阅者的背压不必然阻塞所有订阅者;但如果生产者必须等待最慢订阅者,整体提交仍可能受到影响,具体取决于实现和提交策略。

常见协议错误及其表现

onSubscribe 之前发送数据

错误顺序:

onNext(1)
onSubscribe(subscription)

Subscriber 还没有获得 Subscription,无法表达需求或取消订阅。这破坏了生命周期起点。

没有需求却发送数据

subscriber.onSubscribe(subscription);
// Subscriber 尚未 request
subscriber.onNext(item); // 错误

这会让背压失效。测试中常表现为 Subscriber 收到了它从未请求过的元素。

超过需求量发送

Subscriber 请求一个元素:

subscription.request(1);

Publisher 却连续调用:

subscriber.onNext(1);
subscriber.onNext(2);

第二个元素超出了许可额度。

重复发送终止信号

subscriber.onComplete();
subscriber.onComplete();

或者:

subscriber.onError(ex);
subscriber.onComplete();

终止信号必须是一次性的。

终止后继续发送

subscriber.onComplete();
subscriber.onNext(item);

这会让下游无法判断数据流是否真的结束。

并发调用同一 Subscriber

两个线程同时执行:

subscriber.onNext(a);
subscriber.onNext(b);

即使元素数量没有超过需求量,也可能违反信号串行性和顺序要求。

request 当作线程池大小

subscription.request(100);

表示允许发送 100 个元素,不表示启动 100 个线程,也不表示 Subscriber 会并行处理 100 个元素。并行度是执行模型问题,需求量是流量控制问题,两者不能混为一谈。

如何诊断一个响应式流实现

遇到重复数据、卡住、内存增长或订阅不结束时,可以沿着一条订阅记录以下事件:

订阅建立
onSubscribe
request(n)
onNext(item)
request(n)
cancel / onError / onComplete

至少应记录:

  • 订阅 ID;
  • 当前线程;
  • 每次 request 的参数;
  • 累计请求量;
  • 累计发送量;
  • 缓冲区长度;
  • 是否取消;
  • 是否已经终止。

最有价值的不变量是:

累计发送量 ≤ 累计请求量

如果发现:

累计发送量 > 累计请求量

问题通常位于 Publisher、Processor 或中间操作符,而不是消费业务本身。

如果流“没有任何输出”,应先检查:

  1. onSubscribe 是否被调用;
  2. Subscriber 是否调用了 request(n)
  3. n 是否为正数;
  4. Publisher 是否仍有数据;
  5. 是否被提前 cancel
  6. 是否被异常终止;
  7. 异步执行器是否已经关闭;
  8. Publisher 是否因为缓冲区满而阻塞提交线程。

如果内存持续增长,应检查:

  • Subscriber 是否长期不请求;
  • Publisher 是否有无界缓存;
  • Processor 是否向上游请求过多;
  • 下游需求是否被错误地放大;
  • 生产者是否无法被暂停;
  • 是否存在已经取消但仍保留数据和回调引用的订阅。

规范保证、实现行为和工程选择

需要把三类结论分开。

规范保证

Flow 响应式流协议要求关注:

  • onSubscribe 先于其他信号;
  • request(n) 使用正数;
  • onNext 不得超过需求量;
  • 信号有序且终止信号唯一;
  • 取消后应停止继续发送。

这些是协议正确性的基础。

常见实现行为

具体实现可能选择:

  • 同步调用还是异步调用 onNext
  • 使用线程池、虚拟线程还是事件循环;
  • 需求量使用锁还是 CAS;
  • 缓冲区满时阻塞、失败还是丢弃;
  • Subscriber 异常时如何清理;
  • 多个订阅者之间是否相互影响。

这些不能仅由 Flow.Publisher 这个接口名称推断,必须查看具体实现的 API 文档和源代码。

工程语义选择

以下问题没有统一答案:

  • 慢消费者应该阻塞生产者还是丢数据;
  • 是否需要保序;
  • 是否允许重复处理;
  • 失败后是否重试;
  • 是否需要持久化;
  • 是否需要至少一次或至多一次处理;
  • 是否可以接受无界需求。

Flow 只提供传输协议,不提供完整的消息交付语义。若系统需要确认、重放、持久化和跨进程容错,还需要消息队列、数据库或更高层的响应式库配合。

结语

PublisherSubscriber 描述的是数据流的两端,Subscription 则是控制这条数据流的状态通道。背压的本质是:消费者先声明需求,生产者只能在需求额度内发送

一个实现是否“异步”并不能说明它是否正确。真正需要验证的是:

是否先订阅并建立 Subscription?
是否正确处理 request(n)?
onNext 是否超出需求?
信号是否串行、有序?
终止信号是否唯一?
cancel 是否能停止后续发送?
缓冲和上游请求是否有界?

只有同时满足这些条件,Publisher、Subscriber、背压和错误终止才构成一个协议正确的响应式流。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。