尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RxJava 如何用 generate() 把阻塞式数据源包装成支持背压的 Flowable 并自动清理资源?

RxJava 如何用 generate() 把阻塞式数据源包装成支持背压的 Flowable 并自动清理资源? RxJava 如何用 generate() 把阻塞式数据源包装成支持背压的 Flowable 并自动清理资源【免费下载链接】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如果你的数据源本身是同步、阻塞、拉取式的——必须调用一次read()或get()才能拿到下一个数据——直接把它改成Iterable再交给Flowable.fromIterable()会丢一件重要的事资源。当下游在序列结束前退订时Iterable路径没有地方去关闭流、连接这类资源而阻塞式源往往恰好关联着它们。RxJava 的Flowable.generate()工厂方法族就是为此设计的它产出的Flowable尊重下游背压并且提供第三个回调在序列终止或下游取消时释放资源。本文的任务是把一个读到文件尾才结束的阻塞式数据源如FileInputStream包装成支持背压的Flowable并保证每次订阅独立持有资源、在终止或取消时自动关闭。适用前提是项目里已引入 RxJava 4 的Flowable组件位于io.reactivex.rxjava4包下。准备引入 RxJava 4 依赖README.md 给出的 Gradle 依赖声明把x和y替换为 Maven Central 上最新的版本数字implementation io.reactivex.rxjava4:rxjava:4.x.yREADME 同时说明RxJava 4 的组件位于io.reactivex.rxjava4包基础类与接口位于io.reactivex.rxjava4.core本文代码都按这个包名导入。generate() 的三个回调如何工作docs/Backpressure-(2.0).md.md) 中 Creating backpressured datasources 一节把工厂方法分成两类冷生成器按下游需求返回值热推送器为非反应式源叠加背压处理。generate()属于前者。它接受三个回调Flowable.java 的 Javadoc 也逐条描述了这三者的职责initialStateSupplierS为每个订阅者创建一份独立状态。例子里是new FileInputStream(data.bin)——每个订阅者各自独立打开一次文件而不是共享同一个流。generatorBiFunctionS, EmitterT, S下游每请求值它就执行一次。回调拿到当前状态和一个EmitterT通过onNext/onError/onComplete发信号并返回供下次调用使用的新状态。disposeStateConsumer? super S当下游退订cancel或上面的回调发出终止信号时被调用用来释放资源——例子里关闭FileInputStream。文档特别指出并非所有源都需要全部三个回调Flowable.generate提供了不带某些回调的重载按需选择即可。Javadoc 给出的使用约束同样关键传给回调的Emitter的onNext、onError、onComplete必须同步调用且只能在该次回调函数体执行期间调用不能并发、不能从多个线程调用否则行为未定义同一次回调中多次调用onNext会触发IllegalStateException。所有generate重载都标注了BackpressureSupport(BackpressureKind.FULL)即该算子尊重下游背压。主路径把 FileInputStream 包装成背压感知的 Flowable下面的代码取自 docs/Backpressure-(2.0).md.md) 的generate()一节包名按本仓库 4.x 调整为io.reactivex.rxjava4文档原例使用io.reactiveximport io.reactivex.rxjava4.core.Flowable; import io.reactivex.rxjava4.core.Emitter; import io.reactivex.rxjava4.plugins.RxJavaPlugins; import java.io.FileInputStream; import java.io.IOException; FlowableInteger o Flowable.generate( () - new FileInputStream(data.bin), (inputstream, output) - { try { int abyte inputstream.read(); if (abyte 0) { output.onComplete(); } else { output.onNext(abyte); } } catch (IOException ex) { output.onError(ex); } return inputstream; }, inputstream - { try { inputstream.close(); } catch (IOException ex) { RxJavaPlugins.onError(ex); } } );三个回调在这里的分工() - new FileInputStream(data.bin)把data.bin换成你要读的文件路径。每个订阅者订阅时都会独立执行这句各持有一个流。第二个回调里inputstream.read()是阻塞调用读到负值表示文件结束调用output.onComplete()读出一个字节则output.onNext(abyte)。注意每次调用至多发一次onNext。第三个回调在下游退订或序列终止时执行close()如果close()本身抛出IOException文档示例把它交给RxJavaPlugins.onError(ex)处理避免在清理路径上直接抛出。文档还提醒了一个 Java 语言层面的麻烦generate的函数式接口不允许抛受检异常而 JVM 与各种库的方法调用大量抛受检异常所以这些调用必须包进try-catch——上面的read()和close()都是如此处理的。验证确认背压生效、终止与清理行为仓库测试 FlowableGenerateTest.java 展示了TestSubscriber的几种可复用验证方式1. 有限消费 断言结果。对无终止的源用take(n)截断后断言。FlowableGenerateTest#statefulBiconsumer的做法Flowable.generate(() - 10, (BiConsumerObject, EmitterObject) (s, e) - e.onNext(s), _ - { }) .take(5) .test() .assertResult(10, 10, 10, 10, 10);2. 限制请求量验证不超发。#backpressure测试用.test(5L)把初始请求限制为 5断言恰好收到 5 个值且没有完成——这就是背压生效的表现上游不会在下游请求之外多产数据Flowable.generate(() - 1, (BiConsumerObject, EmitterObject) (_, e) - e.onNext(1), _ - { }) .rebatchRequests(1) .to(TestHelper.ObjecttestSubscriber(5L)) .assertSubscribed() .assertValues(1, 1, 1, 1, 1) .assertNoErrors() .assertNotComplete();把这套模式套到文件读取示例上就是构造TestSubscriber只请求有限个值、中途dispose()然后确认 dispose 回调关闭流被执行、且没有发出超出请求数量的值。3. 用文档的无限序列例子做最简冒烟。docs/Backpressure-(2.0).md.md) 用generate模仿一个无界 rangeFlowable.generate( () - 0, (current, output) - { output.onNext(current); return current 1; }, e - { } );按文档的解释current第一次从0开始lambda 下次被调用时参数current变为1。对这个源直接订阅会无限打印因此文档的类似场景是配合take(5)之类的算子只消费前几个值再停止请求。4. 异常路径。同一测试文件还覆盖了initialState抛异常时订阅者收到onErrorgenerator抛异常时同样收到onErrordisposeState抛出的异常通过插件错误通道报告TestHelper.trackPluginErrors()/assertUndeliverable。如果文件读取在read()阶段抛IOException按主路径代码它会从output.onError(ex)进入订阅者的onError。常见错误点与限制一次回调里调用了两次onNextJavadoc 明确这会发出IllegalStateExceptionFlowableGenerateTest#multipleOnNext断言了assertFailure(IllegalStateException.class, 1)第一个onNext之后的信号被忽略。在回调外或并发调用Emitter方法Javadoc 明确不支持leads to an undefined behavior。阻塞式 I/O 就写在这次回调体内同步执行不要另开线程回调onNext。受检异常函数式接口不允许抛受检异常所有可能抛受检异常的调用都要try-catch后转成output.onError(ex)。类型可用性docs/Creating-Observables.md 的generate一节说明它可用于Flowable和Observable不可用于Maybe、Single、Completable。需要背压时选Flowable.generateObservable.generate产出的Observable本身不支持背压README 对Observable的定位即no backpressure。清理回调里的异常主路径示例没有把close()的异常向上抛而是交给RxJavaPlugins.onError测试文件#disposerThrows也表明 dispose 回调抛出的异常走插件通道报告而不是中断订阅者。文档版本差异docs/Backpressure-(2.0).md.md) 标题带 (2.0)其中create(emitter)等描述针对的是 2.x API本文引用的generate()一节与当前 4.x 源码的 Javadoc 一致Flowable.create的BackpressureStrategy用法请以 Flowable.java 的 Javadoc 为准。完成上面四步验证后你就得到了一个可直接使用的封装每个订阅者独立的阻塞式数据源、按下游请求节奏拉取数据背压感知、在终止或退订时自动清理资源。下一步如果要处理更复杂的取消逻辑比如多个资源参考同一文档中create(emitter)一节介绍的多资源组合方式即可。【免费下载链接】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),仅供参考
返回列表