RxJava 2是否有类似Project Reactor Flux.create()的推拉模型实现?
RxJava 2 中推拉混合模型的实现方案
首先明确:RxJava 2 并没有提供像 Project Reactor 那样专门的工厂方法直接创建支持推拉混合模型的 Producer
核心思路:利用 Flowable.create() 结合背压机制实现推拉
RxJava 2 的 Flowable 是原生支持 Reactive Streams 背压的核心类型,而推拉混合模型的本质就是下游根据自身能力“拉”取数据,上游在收到请求后“推”送对应数据。我们可以通过 Flowable.create() 来包装你的回调式API,结合请求监听逻辑实现这一模型。
示例代码:包装基于回调的条目获取API
假设你有这样一个基于回调的API,每次调用会异步返回一个条目:
// 你现有的回调式API interface ItemFetcher { // 异步获取下一个条目,拿到后调用onItem,无更多数据时调用onComplete void fetchNextItem(Consumer<Item> onItem, Runnable onComplete); }
下面是用 Flowable.create() 包装成推拉模型组件的实现:
import io.reactivex.Flowable; import io.reactivex.FlowableEmitter; import io.reactivex.FlowableOnSubscribe; import java.util.concurrent.atomic.AtomicBoolean; Flowable<Item> createPushPullStream(ItemFetcher fetcher) { return Flowable.create(new FlowableOnSubscribe<Item>() { @Override public void subscribe(FlowableEmitter<Item> emitter) throws Exception { // 防止并发重复请求的标记 AtomicBoolean isFetching = new AtomicBoolean(false); // 封装获取下一个条目的逻辑 Runnable fetchNext = () -> { // 检查是否已取消,且当前没有在获取数据 if (!emitter.isDisposed() && isFetching.compareAndSet(false, true)) { fetcher.fetchNextItem( // 拿到条目后的处理:推送给下游,然后检查是否还有未处理的请求 item -> { isFetching.set(false); emitter.onNext(item); // 如果下游还有未满足的请求,继续获取下一个 if (emitter.requested() > 0) { fetchNext.run(); } }, // 无更多数据时完成流 () -> { isFetching.set(false); emitter.onComplete(); } ); } }; // 设置取消回调:如果下游取消订阅,可在这里清理资源 emitter.setCancellable(() -> { // 示例:如果你的ItemFetcher支持取消,这里可以调用取消方法 // fetcher.cancelCurrentFetch(); }); // 初始阶段:如果下游已有请求,立即开始获取 if (emitter.requested() > 0) { fetchNext.run(); } // 监听下游的请求变化:每当下游调用request(n),就触发获取逻辑 emitter.setOnRequest(n -> { fetchNext.run(); }); } }, BackpressureStrategy.BUFFER); // 根据业务场景选择背压策略,比如BUFFER/DROP/LATEST }
代码说明
- 拉取逻辑:通过
emitter.setOnRequest()监听下游的请求信号,只有当下游主动请求数据时,才会触发fetchNext去调用你的回调API获取条目。 - 推送逻辑:当回调API返回条目后,通过
emitter.onNext(item)将数据推送给下游。 - 背压控制:通过指定
BackpressureStrategy来处理下游请求速度跟不上上游推送的情况,比如BUFFER会缓存未处理的数据,DROP会丢弃超出下游能力的最新数据,可根据业务需求选择。
其他可选方案
如果你的API是同步获取单条数据的,也可以使用 Flowable.generate(),它本身就是基于拉取的模型,每次下游请求时生成一个数据,示例如下:
Flowable<Item> createPullBasedStream(Supplier<Item> syncFetcher) { return Flowable.generate(emitter -> { Item item = syncFetcher.get(); if (item != null) { emitter.onNext(item); } else { emitter.onComplete(); } }); }
不过这种方式更偏向纯拉取,如果你需要结合异步回调的推送逻辑,Flowable.create() 会更灵活。
内容的提问来源于stack exchange,提问作者artur
相关产品推荐
相关产品推荐

