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

Observable.create与fromPublisher的区别是什么?替换后搭配sample无输出如何解决

Observable.create 与 Observable.fromPublisher 的核心区别

两者的设计定位和底层逻辑完全不同:

  • Observable.create是RxJava官方提供的安全自定义Observable API,你拿到的ObservableEmitter已经由RxJava做了全场景封装:自动处理背压适配、订阅生命周期管理、下游请求同步、线程安全校验,你只需要按照逻辑调用onNext/onError/onComplete即可,不需要关心Reactive Streams的底层规范细节。
  • Observable.fromPublisher是Reactive Streams标准适配API,仅用于对接其他符合Reactive Streams规范的Publisher实现,要求你传入的Publisher必须严格遵守规范:必须收到下游的request(long n)请求后才能发送onNext事件,未收到请求时发送的所有事件都会被直接丢弃,同时还要自行处理cancel、线程安全等要求。
sample操作符下fromPublisher失效的原因

sample属于带流量控制的操作符:它订阅上游时不会直接请求无限量数据,而是根据采样周期动态申请事件量。
如果你仿照create的写法直接在Publisher里启动线程就发事件,完全不监听下游的request请求,所有事件都会因为不符合Reactive Streams规范被直接丢弃,自然收不到任何回调。
而不加sample时,普通subscribe默认会在订阅时直接申请Long.MAX_VALUE量的事件,所以你主动发的事件不会被拦截,看起来就像运行正常。

正确的fromPublisher实现参考

如果一定要用fromPublisher实现相同逻辑,需要自行处理request请求:

import io.reactivex.Observable;
import org.reactivestreams.Publisher;
import org.reactivestreams.Subscriber;
import org.reactivestreams.Subscription;

import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;

public class SampleMain {
    public static void main(String[] args) throws InterruptedException {
        Observable<String> o = Observable.fromPublisher(new Publisher<String>() {
            private final AtomicLong requested = new AtomicLong(0);
            @Override
            public void subscribe(Subscriber<? super String> subscriber) {
                subscriber.onSubscribe(new Subscription() {
                    @Override
                    public void request(long n) {
                        requested.addAndGet(n);
                    }

                    @Override
                    public void cancel() {
                        requested.set(-1);
                    }
                });

                new Thread(() -> {
                    for (int i = 0; i < 100; i++) {
                        // 等待下游申请事件
                        while (requested.get() <= 0 && requested.get() != -1) {
                            try {
                                Thread.sleep(1);
                            } catch (InterruptedException e) {
                                e.printStackTrace();
                            }
                        }
                        if (requested.get() == -1) {
                            break;
                        }
                        subscriber.onNext("Hello Observable.fromPublisher() A" + i);
                        subscriber.onNext("Hello Observable.fromPublisher() B" + i);
                        requested.addAndGet(-2);
                        try {
                            Thread.sleep(100);
                        } catch (InterruptedException e) {
                            e.printStackTrace();
                        }
                    }
                    subscriber.onComplete();
                }).start();
            }
        });
        o.sample(1, TimeUnit.MILLISECONDS)
                .subscribe(System.out::println);
        Thread.sleep(11000);
    }
}

日常开发自定义事件源优先用Observable.create即可,不需要手动实现Publisher,后者手动实现很容易违反Reactive Streams规范导致异常。

内容的提问来源于stack exchange,提问作者user212942

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 12:09:04