【问题标题】:How to delay items, but only once in the beginning?如何延迟项目,但在开始时只有一次?
【发布时间】:2018-01-19 03:28:22
【问题描述】:

delay 运算符将所有项目延迟指定的时间量。我只想延迟和缓冲前 N 秒的项目。 N 秒后应该没有延迟。我需要在下面的代码中做到这一点。

private Emitter<Work> workEmitter;

// In the constructor.
Flowable.create(
        (FlowableOnSubscribe<Work>) emitter -> workEmitter = emitter.serialize(),
        BackpressureStrategy.BUFFER)
    .observeOn(Schedulers.from(executor))
    .subscribe(work -> process(work));

// On another thread, as work comes in, ...
workEmitter.onNext(t);

我想要做的是在前 N 秒内推迟处理工作,但在那之后不会。我尝试了延迟订阅,但它在延迟期间将workEmitter 保留为null。我这样做的原因是为了让 CPU 在初始阶段可用于其他重要工作。

【问题讨论】:

    标签: rx-java2


    【解决方案1】:

    您可以使用UnicastProcessor 并在延迟后订阅它:

    FlowableProcessor<Work> processor = UnicastProcessor.<Work>create().toSerialized();
    
    processor.delaySubscription(N, TimeUnit.SECONDS)
    .observeOn(Schedulers.from(executor))
    .subscribe( work -> process(work));
    
    // On another thread, as work comes in, ...
    processor.onNext(t);
    

    UnicastProcessor 将继续缓冲工作项,直到 delaySubscription 的时间过去,然后切换到它。

    【讨论】:

      【解决方案2】:

      您可以延迟创建 observable 然后订阅它。

      Observable.timer( N, SECONDS )
        .flatMap( ignored -> Flowable.create(
          (FlowableOnSubscribe<Work>) emitter -> workEmitter = emitter.serialize(),
             BackpressureStrategy.BUFFER)
          .observeOn(Schedulers.from(executor)))
        .subscribe( work -> process(work));
      

      在 N 秒过去之前,这不会启动观察者链。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2017-11-22
        • 1970-01-01
        • 1970-01-01
        • 2021-11-29
        • 1970-01-01
        • 2022-01-14
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多