RxJava 如何用 publish()、replay() 与 refCount() 让多个订阅者共享同一个数据源?

发布时间:2026/9/10 20:45:45
RxJava 如何用 publish()、replay() 与 refCount() 让多个订阅者共享同一个数据源? RxJava 如何用 publish()、replay() 与 refCount() 让多个订阅者共享同一个数据源【免费下载链接】RxJavaRxJava – Reactive Extensions for the JVM – a library for composing asynchronous and event-based programs using observable sequences for the Java VM.项目地址: https://gitcode.com/gh_mirrors/rx/RxJava当你有两个或更多订阅者都订阅同一个Observable时默认情况下每个订阅者各自触发一次源序列各自收到不同的一份数据。RxJava 的Connectable Observable机制由publish()、replay()产生的ConnectableObservable配合connect()和refCount()解决的正是这个问题让所有订阅者共享同一次发射。适用前提项目已引入 RxJava 4 依赖代码使用io.reactivex.rxjava4.core包下的类型。准备条件引入依赖按 README.md 的说明在 Gradle 中添加implementation io.reactivex.rxjava4:rxjava:4.x.y文档要求把4.x.y中的x、y替换为 Maven Central 上的最新版本号。Java 代码中导入import io.reactivex.rxjava4.core.*;RxJava 4 的组件位于io.reactivex.rxjava4包下基础类和接口位于io.reactivex.rxjava4.core包下。先看清问题普通 Observable 无法共享数据源docs/Connectable-Observable-Operators.md 给出了对照实验同一个rangesample序列两个订阅者先后订阅。情形一订阅普通 ObservableObservable firstMillion Observable.range(1, 1000000).sample(7, java.util.concurrent.TimeUnit.MILLISECONDS); firstMillion.subscribe(next - System.out.println(Subscriber #1: next), // onNext throwable - System.out.println(Error: throwable), // onError () - System.out.println(Sequence #1 complete) // onComplete ); firstMillion.subscribe(next - System.out.println(Subscriber #2: next), // onNext throwable - System.out.println(Error: throwable), // onError () - System.out.println(Sequence #2 complete) // onComplete );文档示例的输出两个订阅者拿到的是两份独立采样的数据数值互不相同Subscriber #1:211128 Subscriber #1:411633 Subscriber #1:629605 Subscriber #1:841903 Sequence #1 complete Subscriber #2:244776 Subscriber #2:431416 Subscriber #2:621647 Subscriber #2:826996 Sequence #2 complete注意上面数值是文档记录的示例输出运行环境不同时会得到不同的采样值但两个订阅者数值不一致这一点是可复现的判断依据。主路径publish() 产生 ConnectableObservableconnect() 统一开始发射Connectable Observable与普通Observable的区别在于订阅时不会开始发射只有调用connect()时才开始。利用这一点可以等所有目标订阅者都订阅之后再统一开始发射。情形二同一代码改为publish()后在两个订阅之后调用connect()ConnectableObservable firstMillion Observable.range(1, 1000000).sample(7, java.util.concurrent.TimeUnit.MILLISECONDS).publish(); firstMillion.subscribe(next - System.out.println(Subscriber #1: next), // onNext throwable - System.out.println(Error: throwable), // onError () - System.out.println(Sequence #1 complete) // onComplete ); firstMillion.subscribe(next - System.out.println(Subscriber #2: next), // onNext throwable - System.out.println(Error: throwable), // onError () - System.out.println(Sequence #2 complete) // onComplete ); firstMillion.connect();文档示例的输出两个订阅者收到的数值完全一致完成事件在全部数据之后Subscriber #2:208683 Subscriber #1:208683 Subscriber #2:432509 Subscriber #1:432509 Subscriber #2:644270 Subscriber #1:644270 Subscriber #2:887885 Subscriber #1:887885 Sequence #2 complete Sequence #1 complete同样这是文档示例输出具体数值会随运行环境变化判断标准是两个订阅者逐条收到相同的数值而不是像情形一那样各收一份。replay()让晚订阅者也能看到完整序列docs/Connectable-Observable-Operators.md 对replay()的说明是确保所有 Subscribers 看到相同的发射序列即使它们在 Observable 开始发射之后才订阅。也就是说publish()场景下晚订阅的订阅者会错过connect()之后已经发出的项换成replay()时晚订阅者会先收到之前已发射的项再跟随后续发射。用 refCount() 免去手动 connect()手动判断所有订阅者都就位了再调用connect()在实际业务里往往不方便。refCount()的作用是让 Connectable Observable 表现得像普通 Observable由引用计数自动决定何时开始发射对应源码中share()就是publish().refCount()的别名见 Observable.java。docs/Backpressure.md 展示了一个实际用法——同一个 bursty 流既要作为数据源、又要作为buffer的关闭信号选择器必须先把源多路复用multicast// we have to multicast the original bursty Observable so we can use it // both as our source and as the source for our buffer closing selector: ObservableInteger burstyMulticast bursty.publish().refCount();这个片段说明publish().refCount()组合的典型场景同一个源序列被多个下游各自消费但不希望消费多次源。验证与一个必须知道的边界验证方式就是文档的对照实验打印两个订阅者收到的项。普通 Observable 下两组数值不一致各收一份publish()connect()或publish().refCount()下两组数值逐条一致即说明数据源已被共享。一个影响行为判断的边界来自 docs/Backpressure.md冷 Observable 被 multicast转换为ConnectableObservable并调用connect()之后实际上就变成了热 Observable在背压和流控意义上应按热 Observable 对待。也就是说共享之后订阅者不再控制源的开始时机这一点在把多路复用用在需要背压控制的流程里时要纳入考虑。相关的操作符清单可在 Alphabetical-List-of-Observable-Operators.md 中查阅connect()、multicast()、publish()、publishLast()、refCount()、replay()均在其中列出。【免费下载链接】RxJavaRxJava – Reactive Extensions for the JVM – a library for composing asynchronous and event-based programs using observable sequences for the Java VM.项目地址: https://gitcode.com/gh_mirrors/rx/RxJava创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询