【问题标题】:How to use RxJava ReplaySubject<object> to emit object in define interval如何使用 RxJava ReplaySubject<object> 在定义间隔内发射对象
【发布时间】:2018-02-28 07:28:55
【问题描述】:

我有一组自定义对象,我想在每 2 秒后发出一次。这意味着当我向 ReplaySubject 添加一组项目时,它应该一个接一个地访问这些对象 @987654321 @ 方法每 2 秒一次。此对象在用户事件上动态添加到 ReplaySubject。这就是我使用 ReplaySubject 的原因。

我已经做到了,但我想暂停并恢复这个线程 基于某些条件,该条件将根据用户交互而改变。如何做到这一点。

这是我的代码。

private ReplaySubject<CartItemToBeRemove> cartItemToBeRemoveSubject;

    @Override
    protected void onCreate(Bundle savedInstanceState) {
        super.onCreate(savedInstanceState);
        setContentView(R.layout.activity_cart);
        ButterKnife.bind(this);

        cartItemToBeRemoveSubject = ReplaySubject.create();

        initDeleteQueue();
    }

public void initDeleteQueue() {

        cartItemToBeRemoveSubject
                .delay(2,TimeUnit.SECONDS)
                .subscribeOn(Schedulers.io())
                .observeOn(AndroidSchedulers.mainThread())
                .subscribeWith(new Observer<CartItemToBeRemove>() {
                    @Override
                    public void onSubscribe(Disposable d) {

                    }

                    @Override
                    public void onNext(CartItemToBeRemove itemToBeRemove) {
                        cartData.remove(itemToBeRemove.getItem());
                        cartListAdapter.notifyItemRemoved(itemToBeRemove.getPosition());
                        Log.d(TAG, "item" + deletedItem.getName());
                    }

                    @Override
                    public void onError(Throwable e) {

                    }

                    @Override
                    public void onComplete() {

                    }
                });
    }

我正在向 ReplaySubject 添加对象,如下所示

cartItemToBeRemoveSubject.onNext(toBeRemove);

【问题讨论】:

  • 我没有明确的答案给你,但你可以看看zipping 你的主题和Observable.interval。

标签: java android rx-java rx-android


【解决方案1】:

您想为可观察流添加节奏。正常节奏是每两秒发出一次项目,但可以暂停和恢复。

基本的起搏机制类似于Observable.interval( 2, SECONDS ),它每2 秒发出一个Long。但是,您希望能够偶尔暂停它。假设我们有一个Observable&lt;Boolean&gt;,它在调步机制要继续时发出TRUE,而在它应该暂停时发出FALSE。

Observable<Boolean> pacingControl;

pacingControl
  .switchMap( ctl -> ctl ? Observable.interval( 0, 2, SECONDS ) : Observable.never() )
  .zipWith( cartItemToBeRemoveSubject, (t, item) -> item )
  .subscribeOn(Schedulers.io())
  .observeOn(AndroidSchedulers.mainThread())
  .subscribe( ... );

pacingControl 可以实现为主题,用户操作将是pacingControl.onNext( TRUE ) 启动进程或pacingControl.onNext( FALSE ) 暂停它。 switchMap() 操作符在定时器和never() 之间切换,这显然不会产生任何结果。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-08-01
    • 1970-01-01
    • 2021-09-11
    • 1970-01-01
    相关资源
    最近更新 更多