You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RxJava 2是否有类似Project Reactor Flux.create()的推拉模型实现?

RxJava 2 中推拉混合模型的实现方案

首先明确:RxJava 2 并没有提供像 Project Reactor 那样专门的工厂方法直接创建支持推拉混合模型的 Producer,但我们完全可以基于现有API在不手动实现响应式规范接口的前提下,快速构建这类组件。

核心思路:利用 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 06:54:09