RX-JAVA distinct debounce defer merge

网友投稿 241 2022-11-24

RX-JAVA distinct debounce defer merge

distinct 去重 Observable.just(1, 1, 1, 2, 2, 3, 4, 5) .distinct() .subscribe(new Consumer() { @Override public void accept(Integer integer) throws Exception { log.info("distinct : " + integer); } }); Thread.sleep(300000000); debounce 去除发送频率过快的项,2个发送项的间隔要大于500毫秒,比如1发完后400毫秒发送2,那么1不会被发送到下游,  2发送完后505毫秒发送3,那么2会发送到下游。 Observable.create(new ObservableOnSubscribe() { @Override public void subscribe(ObservableEmitter emitter) throws Exception { // send events with simulated time wait emitter.onNext(1); // skip Thread.sleep(400); emitter.onNext(2); // deliver Thread.sleep(505); emitter.onNext(3); // skip Thread.sleep(100); emitter.onNext(4); // deliver Thread.sleep(605); emitter.onNext(5); // deliver Thread.sleep(510); emitter.onComplete(); } }).debounce(500, TimeUnit.MILLISECONDS) .subscribeOn(Schedulers.io()) .observeOn(Schedulers.newThread()) .subscribe(new Consumer() { @Override public void accept(Integer integer) throws Exception { log.info("debounce :" + integer); } }); Thread.sleep(300000000);   defer 简单地时候就是每次订阅都会创建一个新的 Observable,并且如果没有被订阅,就不会产生新的 Observable。 Observable observable = Observable.defer(new Callable>() { @Override public ObservableSource call() throws Exception { log.info("call -create observable"); return Observable.just(1, 2, 3); } }); observable.subscribe(new Consumer() { @Override public void accept(Integer integer) throws Exception { log.info("accpent:" + integer); } }); observable.subscribe(new Consumer() { @Override public void accept(Integer integer) throws Exception { log.info("accpent:" + integer); } }); 2021-02-24 10:53:24 [main] INFO c.n.d.v.RxJavaTest:call - call -create observable 2021-02-24 10:53:24 [main] INFO c.n.d.v.RxJavaTest:accept - accpent:1 2021-02-24 10:53:24 [main] INFO c.n.d.v.RxJavaTest:accept - accpent:2 2021-02-24 10:53:24 [main] INFO c.n.d.v.RxJavaTest:accept - accpent:3 2021-02-24 10:53:24 [main] INFO c.n.d.v.RxJavaTest:call - call -create observable 2021-02-24 10:53:24 [main] INFO c.n.d.v.RxJavaTest:accept - accpent:1 2021-02-24 10:53:24 [main] INFO c.n.d.v.RxJavaTest:accept - accpent:2 2021-02-24 10:53:24 [main] INFO c.n.d.v.RxJavaTest:accept - accpent:3 merge 顾名思义,熟悉版本控制工具的你一定不会不知道 merge 命令,而在 Rx 操作符中,merge 的作用是把多个 Observable 结合起来,接受可变参数,也支持迭代器集合。注意它和 concat 的区别在于,不用等到 发射器 A 发送完所有的事件再进行发射器 B 的发送。 Observable.merge(Observable.just(1, 2, 3, 4), Observable.just(5, 6, 7, 8)) .subscribeOn(Schedulers.newThread()) .observeOn(Schedulers.io()) .subscribe(new Consumer() { @Override public void accept(Integer integer) throws Exception { log.info("accept: merge :" + integer); } }); Thread.sleep(3000000);

版权声明:本文内容由网络用户投稿,版权归原作者所有,本站不拥有其著作权,亦不承担相应法律责任。如果您发现本站中有涉嫌抄袭或描述失实的内容,请联系我们jiasou666@gmail.com 处理,核实后本网站将在24小时内删除侵权内容。

上一篇:关于Java 项目封装sqlite连接池操作持久化数据的方法
下一篇:Ali-HBase的SQL实践与改进
相关文章

 发表评论

暂时没有评论,来抢沙发吧~