【问题标题】:Creating Multiple Timers with Reactive Extensions使用响应式扩展创建多个计时器
【发布时间】:2013-08-08 01:26:01
【问题描述】:

我有一个非常简单的类,我用它来轮询目录中的新文件。它有位置、开始监控该位置的时间以及再次检查的时间间隔(以小时为单位):

public class Thing
{
  public string         Name {get; set;}
  public Uri            Uri { get; set;}
  public DateTimeOffset StartTime {get; set;}
  public double         Interval {get; set;}
}

我是 Reactive Extensions 的新手,但我认为它正是适合这里工作的工具。在开始时以及随后的每个间隔,我只想调用一个 Web 服务来完成所有繁重的工作 - 我们将使用富有创意的 public bool DoWork(Uri uri) 来表示它。

编辑: DoWork 是对 Web 服务的调用,它将检查新文件并在必要时移动它们,因此它的执行应该是异步的。如果完成则返回 true,否则返回 false。

如果我有这些Things 的完整集合,事情就会变得复杂。我不知道如何为每个人创建Observable.Timer(),并让他们都调用相同的方法。

edit2: Observable.Timer(DateTimeOffset, Timespan) 似乎非常适合为我在这里尝试做的事情创建一个 IObservable。想法?

【问题讨论】:

    标签: timer system.reactive reactive-programming


    【解决方案1】:

    你需要有很多计时器吗?我假设如果你有 20 个东西的集合,那么我们将创建 20 个计时器来在同一时间点全部触发?在同一个线程/调度器上?

    或者也许你想DoWork foreach 在每个时期的事情?

    from thing in things
    from x in Observable.Interval(thing.Interval)
    select DoWork(thing.Uri)
    

    对比

    Observable.Interval(interval)
    .Select(_=>
        {
            foreach(var thing in Things)
            {
                DoWork(thing);
            }
        })
    

    您可以通过多种方式在未来开展工作。

    • 您可以直接使用调度程序来安排要在 未来。
    • 您可以使用 Observable.Timer 让序列在未来的指定时间产生一个值。
    • 您可以使用 Observable.Interval 来创建一个序列,该序列在每个指定的时间段内产生许多值。

    所以现在引入另一个问题。如果您的轮询时间为 60 秒,而您的工作功能需要 5 秒;下一次投票应该在 55 秒内还是 60 秒内进行?这里一个答案表示您想要使用 Rx 序列,另一个表示您可能想要使用 Periodic Sc​​heudling。

    下一个问题是,DoWork 是否返回值?目前看起来它没有*。在这种情况下,我认为最适合您的做法是利用定期调度程序(假设为 Rx v2)。

    var things = new []{
        new Thing{Name="google", Uri = new Uri("http://google.com"), StartTime=DateTimeOffset.Now.AddSeconds(1), Interval=3},
        new Thing{Name="bing", Uri = new Uri("http://bing.com"), StartTime=DateTimeOffset.Now.AddSeconds(1), Interval=3}
    };
    var scheduler = Scheduler.Default;
    var scheduledWork = new CompositeDisposable();
    
    foreach (var thing in things)
    {
        scheduledWork.Add( scheduler.SchedulePeriodic(thing, TimeSpan.FromSeconds(thing.Interval), t=>DoWork(t.Uri)));
    }
    
    //Just showing that I can cancel this i.e. clean up my resources.
    scheduler.Schedule(TimeSpan.FromSeconds(10), ()=>scheduledWork.Dispose());
    

    现在这将安排每件事情定期处理(没有漂移),在它自己的时间间隔内并提供取消。

    如果您愿意,我们现在可以将其升级为查询

    var scheduledWork = from thing in things
                        select scheduler.SchedulePeriodic(thing, TimeSpan.FromSeconds(thing.Interval), t=>DoWork(t.Uri));
    
    var work = new CompositeDisposable(scheduledWork);
    

    这些查询的问题是我们没有满足StartTime 的要求。令人讨厌的是,Ccheduler.SchedulePeriodic 方法没有提供重载来也有一个起始偏移量。

    Observable.Timer 运算符确实提供了这一点。它还将在内部利用非漂移调度功能。要使用Observable.Timer 重构查询,我们可以执行以下操作。

    var urisToPoll = from thing in things.ToObservable()
                     from _ in Observable.Timer(thing.StartTime, TimeSpan.FromSeconds(thing.Interval))
                     select thing;
    
    var subscription = urisToPoll.Subscribe(t=>DoWork(t.Uri));
    

    所以现在你有一个很好的界面,应该避免漂移。但是,我认为这里的工作是以串行方式完成的(如果同时调用许多 DoWork 动作)。

    *理想情况下,我会尽量避免这样的副作用陈述,但我不能 100% 确定您的要求。

    编辑 看来对 DoWork 的调用必须是并行的,所以你需要做更多的事情。理想情况下,您将 DoWork 设为 asnyc,但如果您无法做到,我们可以在成功之前对其进行伪造。

    var polling = from thing in things.ToObservable()
                  from _ in Observable.Timer(thing.StartTime, TimeSpan.FromSeconds(thing.Interval))
                  from result in Observable.Start(()=>DoWork(thing.Uri))
                  select result;
    
    var subscription = polling.Subscribe(); //Ignore the bool results?
    

    【讨论】:

    • 这是一个很棒的解释——尽管我还在消化它——DoWork 是对 Web 服务的调用。它会检查每个文件夹 (uri) 中的新文件,并对找到的文件执行操作……可怕、可怕的事情。
    • “漂移” - 根据前一个 DoWork 的完成来补偿下一个间隔 - 不是必需的。预计间隔是最频繁的每半小时一次。
    • 您可以忽略所有这些,而只使用最后一个查询/代码块。其他答案也都不错。
    • 我不明白 "from _ in Observable.Timer" 行是如何做任何事情的。你永远不会将结果用于任何事情......
    • 在函数式编程中(至少在 C# 中),下划线字符是礼貌的做法,表示不使用返回值或输入参数。在这种情况下,它不被使用。从语法上讲,那里需要一个变量名。我们需要这条线,因为我们想要为每个“事物”设置一个计时器。 Timer 序列将产生一个值(忽略),然后允许查询链执行 Observable.Start(DoWork) 子句
    【解决方案2】:

    嗯。 DoWork 是否会产生您需要处理的结果?我会假设是这样。你没有说,但我也会假设 DoWork 是同步的。

    things.ToObservable()
      .SelectMany(thing => Observable
        .Timer(thing.StartTime, TimeSpan.FromHours(thing.Interval))
        .Select(_ => new { thing, result = DoWork(thing.Uri) }))
      .Subscribe(x => Console.WriteLine("thing {0} produced result {1}",
                                        x.thing.Name, x.result));
    

    这是一个带有假设 Task<bool> DoWorkAsync(Uri) 的版本:

    things.ToObservable()
      .SelectMany(thing => Observable
        .Timer(thing.StartTime, TimeSpan.FromHours(thing.Interval))
        .SelectMany(Observable.FromAsync(async () =>
           new { thing, result = await DoWorkAsync(thing.Uri) })))
      .Subscribe(x => Console.WriteLine("thing {0} produced result {1}",
                                        x.thing.Name, x.result));
    

    此版本假定DoWorkAsync 将在间隔到期之前很久完成并启动一个新实例,因此不会防止并发DoWorkAsync 为同一个Thing 实例运行。

    【讨论】:

    • DoWork 调用 Web 服务来检查文件夹中的文件。如果成功完成则返回布尔值,否则返回 false。
    • "我假设 DoWork 是同步的" - DoWork 应该在每个 Interval 中为每个事物执行一次,但独立于其他事物。事情应该并行执行... DoWork 调用 Web 服务来检查每个文件夹中的文件。
    • 那么您可能希望在对 DoWork 的调用中引入一些并发性。更好的是,通过返回 Task 或 IObservable 使 DoWork 异步,然后你可以将它们很好地链接在一起。
    • 是的,这个版本是并行运行的,所以如果多个 Things 时间间隔重合,可能会有并发的 DoWork 调用。我同意@LeeCampbell 的建议,如果DoWork 正在调用某个Web 服务,如果它是异步方法,您可以提高资源利用率。
    • @Brandon 我在result = await DoWorkAsync(thing.Uri) 下得到一个红色波浪线,并显示消息“无法将 void 分配给匿名类型属性”。虽然我确信这告诉了我一些有用的东西,但如果我知道它是什么,我会被诅咒的。
    【解决方案3】:

    这是我能想到的最清晰的方法:

    foreach (var thing in things)
    {
        Scheduler.ThreadPool.Schedule(
            thing.StartTime,
            () => {
                    DoWork(thing.Uri); 
                    Observable.Interval(TimeSpan.FromHours(thing.Interval))
                              .Subscribe(_ => DoWork(thing.Uri)));
                  }
    }
    

    下面这个功能比较多:

    things.ToObservable()
        .SelectMany(t => Observable.Timer(t.StartTime).Select(_ => t))
        .SelectMany(t => Observable.Return(t).Concat(
            Observable.Interval(TimeSpan.FromHours(t.Interval)).Select(_ => t)))
        .Subscribe(t => DoWork(t.Uri));
    

    First SelectMany 在其预定时间创建一个事物流。第二个 SelectMany 采用此流并创建一个在每个间隔重复的新流。这需要加入到Observable.Return,它会立即产生一个值,因为Observable.Interval 的第一个值被延迟。

    *注意,第一个解决方案需要 c#5 否则 this 会咬你。

    【讨论】:

    • 你为什么选择这个方法而不是 Observable.Timer(DateTimeOffset, TimeSpan) ?它不会产生正确的“事件”吗?
    • @ScottSEA 我只是忘记了超载:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-11-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-05-17
    • 1970-01-01
    相关资源
    最近更新 更多