【问题标题】:Momentarily ignore values from Observable when another Observable provides a value当另一个 Observable 提供值时,暂时忽略来自 Observable 的值
【发布时间】:2015-06-25 02:23:43
【问题描述】:

当另一个 Observable 提供值时,我需要忽略 Observable 值一段时间

目前,我的实现使用一个变量来控制阻塞(或忽略)。

bool block = false;

var blocker = observable1.Do(_ => block = true )
                         .Throttle( _ => Observable.Timer(_timeToBlock)
                         .Subscribe( _ => block = false ));

var receiver = observable2.Where( i => !block && SomeCondition(i) )
                          .Subscribe( i=> EvenMoreStuff(i) );

通过结合这两个 observables,是否有更多的 Rx 方法来做到这一点?

编辑:对拦截器订阅的小改动

【问题讨论】:

    标签: .net system.reactive


    【解决方案1】:

    第一个任务是将您的block 变量表示为可观察的。

    IObservable<bool> blockObservable = observable1
        .Select(x => Observable.Concat(
            Observable.Return(true),
            Observable.Return(false).Delay(_timeToBlock)))
        .Switch()
        .DistinctUntilChanged();
    

    每次observable1 发出一个值,我们选择一个发出true 的可观察对象,等待_timeToBlock,然后发出falseSwitch 总是切换到最新的 observables。

    这是一个大理石图。假设_timeToBlock 的长度为 3 个字符。

    observable1      -----x--xx-x-----x-----
    select0               T--F
    select1                  T--F
    select2                   T--F
    select3                     T--F
    select4                           T--F
    switch           -----T--TT-T--F--T--F--
    blockObservable  -----T--------F--T--F--
    

    现在我们可以使用blockObservableMostRecent 值来Zip 您的值序列。

    var receiverObservable = observable2
        .Zip(blockObservable.MostRecent(false), (value, block) => new { value, block })
        .Where(x => !x.block)
        .Select(x => x.value);
    

    【讨论】:

    • 看起来不错。当然解决方案是用 Observable&lt;bool&gt; 替换 bool !我需要围绕语法进行一些测试......
    • @supertopi 这个“浏览器代码”是否在没有任何调整的情况下工作?
    • 至少我的单元测试通过了。我应该知道任何隐藏的副作用吗? :)
    • @supertopi 不,应该没问题。
    【解决方案2】:

    作为使用一次性用品的替代方法,您可以创建一个小的扩展方法:

    public static IObservable<TResult> Suppress<TResult, TOther>(
                                        this IObservable<TResult> source, 
                                             IObservable<TOther> other,
                                             TimeSpan delayFor)
    {
      return Observable.Create<TResult>(observer => {
        var published = source.Publish();
        var connected = new SerialDisposable();
        Func<IDisposable> connect = () => published.Subscribe(observer);
        var suppressor = other.Select(_ => Observable.Timer(delayFor)
                                    .Select(_2 => connect())
                                    .StartWith(Disposable.Empty))
        .Switch();
    
        return new CompositeDisposable(
                    connected,
                    suppressor.StartWith(connect())
                              .Subscribe(d => connected.Disposable = d),
                    published.Connect());
      });
    }
    

    这会将可观察的源转换为ConnectableObservable,然后每次other 源发出它都会释放其订阅,然后在计时器到期时重新连接。

    【讨论】:

    • 谢谢。编写自己的扩展方法可能是最佳实践,但我发现这比 Timothy 的解决方案更难阅读。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-03-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多