Java Flow 是 Java 9+ 原生异步响应式数据流方案,基于发布-订阅模型与背压机制;通过 SubmissionPublisher 发布、自定义 Subscriber 主动请求数据、Processor 转换流、配合 Executor 控制线程上下文实现高效可控的流处理。

Java 中使用 Flow 实现异步数据流处理,核心是依托 java.util.concurrent.Flow 接口标准,构建发布-订阅模型,并借助背压机制保障生产与消费速率协调。它不依赖第三方库,是 Java 9+ 原生支持的轻量级响应式方案,适合对依赖精简、线程控制明确的场景。
用 SubmissionPublisher 快速搭建发布者
SubmissionPublisher 是 Flow API 提供的开箱即用发布者实现,内部基于 ForkJoinPool 异步提交,天然支持多线程安全和背压反馈。
- 创建时可指定缓冲区大小(如
new SubmissionPublisher(ForkJoinPool.commonPool(), 16)),该值影响未被消费数据的暂存上限 - 调用
submit()发布元素,若缓冲区满且订阅者未及时请求,会阻塞或抛出InterruptedException(取决于构造参数) - 推荐配合 try-with-resources 使用,确保
close()被调用,自动触发onComplete()
实现 Subscriber 并主动管理背压
订阅者必须自己控制数据拉取节奏,这是 Flow 区别于“推即发”的关键。不能被动接收,而要通过 subscription.request(n) 显式申领。
在 Java 中初始化和管理阿里云 SDK客户端。包括单例模式、线程安全、endpoint 与 region 配置、VPC 终端节点、同步与异步等。
-
onSubscribe()中必须立即调用request(1)或其他正整数,否则不会收到任何数据 -
onNext()处理完一个元素后,通常再调用request(1)继续获取下一个,形成“一来一请”节拍 - 若想批量处理,可在
onSubscribe()一次性request(10),后续在onNext()中计数,待积攒够再统一处理 - 遇到异常应调用
subscription.cancel(),避免资源泄漏;onComplete()不需再 request
用 Processor 桥接并转换数据流
当需要对数据做中间处理(如过滤、格式转换、聚合),又希望保持响应式链路不断,可用 Processor —— 它既是 Subscriber 又是 Publisher。
立即学习“Java免费学习笔记(深入)”;
- 典型做法:继承
SubmissionPublisher并实现Subscriber接口,复用其发布能力 - 在
onNext()中完成转换逻辑(如字符串转大写、JSON 解析),再调用super.submit(transformed)向下游发布 - 注意同步问题:若上游并发提交,
onNext()可能被多线程调用,需对共享状态加锁或使用线程安全容器
结合 Executor 控制执行上下文
Flow API 本身不绑定线程模型,所有回调(onNext 等)默认在发布者线程执行。如需切换线程(比如 I/O 操作放 IO 线程、计算放 CPU 线程),需手动调度。
- 在
onNext()内部用executor.execute(() -> { /* 耗时操作 */ })异步分发 - 若需返回结果给下游,可用
CompletableFuture封装,再通过thenAccept回传给submit() - 避免在
onNext()中直接阻塞(如Thread.sleep),否则会卡住整个发布链路


















