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
相关产品推荐
相关产品推荐

