【问题标题】:Rx to trigger an action after a certain amount of timeRx 在一定时间后触发一个动作
【发布时间】:2014-10-24 06:26:07
【问题描述】:

我有一个类,它有几个 bool 属性,并且订阅了一个 observable,它提供包含各种值的对象(以不确定的速度)。例如:

bool IsActive {get; private set;}
bool IsBroken {get; private set;}
bool Status {get; private set;}
...
valueStream.GetValues().Subscribe(UpdateValues);

UpdateValues 根据传入的对象做一些工作。不过,我有兴趣将一个值用于某些特定的逻辑。我们称它为 obj.SpecialValue:

private void UpdateValues(MyObject obj)
{
  ...
  Status = obj.SpecialValue;
  ...
}

现在,如果 IsActive 设置为 true 并且 Status 为 false,我想在将 IsBroken 设置为 true 之前等待三秒钟,以便让流有机会在该时间返回 true 的 obj.SpecialValue。如果在三秒内它确实返回 true,我什么也不做,否则将 IsBroken 设置为 true。当 Status 或 IsActive 更新时,再次检查。

一开始我有:

Observable.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(3)).Subscribe(SetIsBroken);

private void SetIsBroken()
{
  IsBroken = IsActive && !Status;
}

但这比它需要做的检查更多。它只需要检查流更新或 IsActive 何时更改。

关于如何正确执行此操作的任何提示?

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    使用BehavourSubject<T> 支持属性

    解决这个问题的一个有用的想法是使用BehaviorSubject<bool> 类型支持您的属性。这些有用地服务于 active 作为属性和该属性的值流的双重目的。

    您可以将它们作为 observable 订阅,也可以通过 Value 属性访问它们的当前值。您可以通过OnNext 发送一个新值来更改它们。

    例如,我们可以这样做:

    private BehaviorSubject<bool> _isActive;
    
    public bool IsActive
    {
        get { return _isActive.Value; }
        set { _isActive.OnNext(value); }
    }
    

    在您的所有属性中都设置了这一点后,就可以非常简单地查看属性是否符合您陈述的复杂条件。假设 _status_isBroken 是类似实现的支持主题,我们可以像这样设置订阅:

    Observable.CombineLatest(_isActive,
                             _status,
                            (a,s) => a & !s).DistinctUntilChanged()
              .Where(p => p)
              .SelectMany(_ => Observable.Timer(TimeSpan.FromSeconds(3), scheduler)
                                         .TakeUntil(_status.Where(st => st)))
              .Subscribe(_ => _isBroken.OnNext(true));  
    

    部分行使用CombineLatest 并订阅_isActive_status 流。每当这些变化中的任何一个发生变化时,它都会发出 - 当_isActive 为真且_status 为假时,结果函数恰好设置了一个真值。 DistinctUntilChanged() 可防止将 _isActive_status 设置为它们已有的值以启动新计时器。

    然后我们使用Where 只过滤这个条件。

    SelectMany 现在将采用真值并将每个值投射到 3 秒后发出的流中,使用 Timer - 但是我们使用 TakeUntil 压缩这个值_status 变为真的事件。 SelectMany 还将流流平展回单个布尔流。

    这里不确定 - 你没有提到它,但你可能想考虑 _isActive 变为 false 是否也应该终止计时器。如果是这种情况,您可以使用MergeTakeUntil 中将对此的监视和_status 结合起来。

    我们现在可以订阅这整个事情来设置_isBroken true 如果这个查询被触发,表明计时器已过期。

    注意 Timerscheduler 参数 - 这是存在的,因此我们可以传入测试调度程序。

    我不确定我是否正确地捕捉到了您的所有逻辑 - 但如果不是,希望您能看到如何根据需要进行修改。

    这是完整的示例。使用 nuget 包 rx-testing,这将在 LINQPad 中运行:

    void Main()
    {
        var tests = new Tests();
        tests.Test();
    }
    
    public class Foo
    {
        private BehaviorSubject<bool> _isActive;
        private BehaviorSubject<bool> _isBroken;
        private BehaviorSubject<bool> _status;
    
        public bool IsActive
        {
            get { return _isActive.Value; }
            set { _isActive.OnNext(value); }
        }
    
        public bool IsBroken { get { return _isBroken.Value; } }
        public bool Status { get { return _status.Value; } }
    
        public Foo(IObservable<MyObject> valueStream, IScheduler scheduler)
        {
            _isActive = new BehaviorSubject<bool>(false);
            _isBroken = new BehaviorSubject<bool>(false);
            _status = new BehaviorSubject<bool>(false);
    
            // for debugging purposes
            _isActive.Subscribe(a => Console.WriteLine(
                 "Time: " + scheduler.Now.TimeOfDay + " IsActive: " + a));
            _isBroken.Subscribe(a => Console.WriteLine(
                 "Time: " + scheduler.Now.TimeOfDay + " IsBroken: " + a));
            _status.Subscribe(a => Console.WriteLine(
                  "Time: " + scheduler.Now.TimeOfDay + " Status: " + a));
    
            valueStream.Subscribe(UpdateValues);
    
            Observable.CombineLatest(
                        _isActive,
                        _status,
                        (a,s) => a & !s).DistinctUntilChanged()
                    .Where(p => p)
                    .SelectMany(_ => Observable.Timer(TimeSpan.FromSeconds(3),
                                                      scheduler)
                                               .TakeUntil(_status.Where(st => st)))
                    .Subscribe(_ => _isBroken.OnNext(true));            
        }
    
        private void UpdateValues(MyObject obj)
        {
            _status.OnNext(obj.SpecialValue);
        }   
    }   
    
    public class MyObject
    {
        public MyObject(bool specialValue)
        {
            SpecialValue = specialValue;
        }
    
        public bool SpecialValue { get; set; }
    }
    
    public class Tests : ReactiveTest
    {
        public void Test()
        {
            var testScheduler = new TestScheduler();
    
            var statusStream = testScheduler.CreateColdObservable<bool>(
                OnNext(TimeSpan.FromSeconds(1).Ticks, false),
                OnNext(TimeSpan.FromSeconds(3).Ticks, true),
                OnNext(TimeSpan.FromSeconds(5).Ticks, false));
    
            var activeStream = testScheduler.CreateColdObservable<bool>(
                OnNext(TimeSpan.FromSeconds(1).Ticks, false),
                OnNext(TimeSpan.FromSeconds(6).Ticks, true));           
    
            var foo = new Foo(statusStream.Select(b => new MyObject(b)), testScheduler);
    
            activeStream.Subscribe(b => foo.IsActive = b);       
    
            testScheduler.Start();               
        }        
    }
    

    回复评论

    如果你想让 isActive false 设置 isBroken false,那么我认为这加起来现在是这样说的:

    isActive isStatus Action
    T        F        Set Broken True after 3 seconds unless any other result occurs
    T        T        Set Broken False immediately if not already false, cancel timer
    F        F        Set Broken False immediately if not already false, cancel timer
    F        T        Set Broken False immediately if not already false, cancel timer
    

    在这种情况下,请使用以下查询:

    Observable.CombineLatest(
                _isActive,
                _status,
                (a,s) => a & !s).DistinctUntilChanged()                    
            .Select(p => p ? Observable.Timer(TimeSpan.FromSeconds(3),
                                                   scheduler)
                                       .Select(_ => true)
                           : Observable.Return(false))
            .Switch()
            .DistinctUntilChanged()
            .Subscribe(res => _isBroken.OnNext(res));
    

    注意变化:

    • SelectMany 现在是 Select,它将每个事件变成
      • 3 秒后返回 true 的计时器
      • 或直接false
    • Select 的结果是一个布尔流:IObservable&lt;IObservable&lt;bool&gt;&gt;。我们希望任何出现的新流都取消任何先前的流。这就是Switch 将要做的——在过程中使结果变平。
    • 我们现在应用第二个 DistinctUntilChanged(),因为取消的计时器可能会导致两个错误值连续出现在流中
    • 最后我们将新出现的布尔值赋给isBroken

    【讨论】:

    • 谢谢詹姆斯。这真的很有帮助。关于 _isActive 变为假,在这种情况下,我想立即将 _isBroken 设置为假(无需等待任何计时器)。我试图根据 _isActive 和 _status 值修改您的代码以将 _isBroken 设置为 true 或 false,但我遇到了计时器问题。我不知道如何保留计时器并仍然将单个布尔值传递给订阅(类似于订阅(b => _isBroken.OnNext(b))。计时器返回一个长整数,我试图从中选择一个布尔值有一点我订阅的时候还是很长。我一定是错过了什么。
    • 听起来您需要更改查询...我已对答案进行了编辑。
    猜你喜欢
    • 2015-05-05
    • 1970-01-01
    • 1970-01-01
    • 2020-10-26
    • 1970-01-01
    • 2015-12-08
    • 2016-10-02
    • 2014-09-30
    • 1970-01-01
    相关资源
    最近更新 更多