ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

RxJava Subject深度解析:Publish、Behavior、Replay等5种Subject玩转热数据流

RxJava Subject深度解析:Publish、Behavior、Replay等5种Subject玩转热数据流 RxJava Subject深度解析Publish、Behavior、Replay等5种Subject玩转热数据流【免费下载链接】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在RxJava中Subject是把冷数据流变成热数据流的关键开关——它既是观察者Observer又是数据源Observable能将一个事件源分发给多个订阅者。本文将带你快速掌握PublishSubject、BehaviorSubject、ReplaySubject、AsyncSubject、UnicastSubject这 5 种 Subject 的用法与区别附选型对照表帮你在事件总线、状态共享、聊天消息等热数据流场景中做出正确选择。什么是 Subject冷热流转换的桥梁普通 Observable 是冷的只有被订阅时才产生数据每个订阅者都会触发一次完整的上游执行。而 Subject 本身就是一个活的事件源向上看它实现了Observer接口可以接收onNext/onError/onComplete事件向下看它继承自Observable可以把接收到的事件多播multicast给任意多个订阅者。这就像广播电台主播数据源只说一遍收听的人订阅者可以在任意时刻打开收音机。其核心定义在 Subject.java 中官方概念说明可参考 docs/Subject.md。5种Subject速览一张表看懂区别Subject 类型是否记忆历史数据新订阅者立刻收到典型场景 PublishSubject❌ 无记忆订阅后产生的新事件事件总线、实时通知 BehaviorSubject⚪ 仅最新一个值最新值 后续新事件用户登录态、当前选中项️ ReplaySubject✅ 全部/最近N个/最近T时间缓冲区内的历史事件 新事件聊天历史、日志回放 AsyncSubject⚪ 仅最终值流结束complete后收到最终值收集异步任务的最终结果 UnicastSubject 缓冲全部单订阅者缓冲的所有数据单消费者的内存友好缓冲PublishSubject零记忆只转发当下PublishSubject是最朴素的 Subject它不缓存任何数据只把订阅之后收到的事件转发出去。早来的人错过的事永远不会补发。创建方式工厂方法无公开构造器PublishSubjectString subject PublishSubject.create(); subject.subscribe(value - System.out.println(value)); // 只会收到此行之后的事件✅适用全局事件总线EventBus 风格、推送通知、实时数据转发 ❌注意订阅前的事件会丢失无人订阅时事件直接被丢弃可用hasObservers()检查源码见 PublishSubject.java行为测试参考 PublishSubjectTest.java。BehaviorSubject自动补发最新状态如果希望新订阅者一进来就拿到当前最新值就用BehaviorSubject。它只缓存最近一个元素订阅时先补发该值再转发后续事件。BehaviorSubject.create()创建无初始值的实例BehaviorSubject.createDefault(value)指定默认值订阅时立即收到getValue()/hasValue()非阻塞、线程安全地读取最新缓存值见 BehaviorSubject.java⚠️ 注意Subject 终止complete/error后缓存值会被清空晚到的订阅者只能收到终止信号不允许null作为值。✅适用SharedPreferences 式状态共享、当前城市/主题、登录用户信息ReplaySubject完整回放历史数据ReplaySubject是记忆力最强的 Subject提供 4 种创建策略见 ReplaySubject.java创建方法行为create()无界缓存所有历史事件createWithSize(maxSize)只保留最近 N 个事件createWithTime(maxAge, unit, scheduler)只保留最近 T 时间内的事件createWithTimeAndSize(...)同时按数量 时间双限制最省内存✅适用聊天记录新房间成员拉取最近消息、日志回放、股票最近报价 小建议生产环境优先用带maxSize或maxAge的重载避免无界缓存拖垮内存AsyncSubject只关心最终结果AsyncSubject是另一种极简主义者流进行期间它只默默记住最后一个值直到onComplete()才把该值一次性发给所有订阅者若以错误终止则什么值都不发。✅适用多个异步操作汇聚后的最终汇总值、表单提交的最终校验结果 ❌注意流不结束就收不到任何数据务必保证上游能正常 complete源码见 AsyncSubject.java。UnicastSubject单订阅者的缓冲筒UnicastSubject是 RxJava 2.x 新增的隐藏款它把上游数据全部缓冲到队列中但只允许一个订阅者第二个订阅者会直接抛IllegalStateException。UnicastSubject.create()及带capacityHint/delayError的重载可预分配容量、控制错误延迟策略见 UnicastSubject.java内部用数组队列实现比ReplaySubject的链表缓冲更省内存✅适用明确只有一个消费者的预缓冲场景且下游稍晚才订阅 ❌注意单订阅限制是设计使然别把它当事件总线用如何选型3个判断问题新订阅者需要补发历史吗不需要 →PublishSubject只补最新值 →BehaviorSubject补一段历史 →ReplaySubject只关心最终结果吗是 →AsyncSubject确定只有一个消费者且想省内存吗是 →UnicastSubject使用Subject的3个安全要点线程安全onNext/onError/onComplete必须串行调用。多生产者场景请调用 Subject.java 中的toSerialized()它会自动包装为 SerializedSubject同时保护重入调用状态探测所有 Subject 提供hasObservers()、hasComplete()、hasThrowable()等线程安全的状态查询便于在推送前判断是否还有接收者及时取消订阅者离开页面/组件销毁时务必 dispose避免 Subject 持续持有引用造成内存泄漏进一步阅读官方 Subject 概念文档docs/Subject.md5 种实现源码src/main/java/io/reactivex/rxjava4/subjects/热数据流的背压处理背景docs/Backpressure-(2.0).md.md)掌握这 5 种 Subject你就拥有了 RxJava 热数据流的全部遥控器用 Publish 广播当下、Behavior 同步状态、Replay 回放历史、Async 汇聚结果、Unicast 缓冲投递。选对工具冷热流的切换就会变得简单而优雅 【免费下载链接】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),仅供参考
返回列表