【问题标题】:Observer pattern with timer带计时器的观察者模式
【发布时间】:2011-04-13 11:08:15
【问题描述】:

我在我的应用程序中使用了观察者模式。

我有一个主题,其中有一个名为 'tmr' 的 System.Timers.Timer 对象。此计时器的滴答事件每 60 秒 触发一次。在这个滴答事件中,我将通知所有与我的主题相关的观察者。我使用了一个 for 循环来遍历我的观察者列表,然后触发观察者更新方法。

假设我有 10 个观察者附加到我的主题。

每个观察者需要 10 秒来完成其处理。

现在在 for 循环中完成通知会导致最后一个观察者的更新方法在 90 秒后被调用。即下一个观察者更新方法仅在前一个完成处理后调用。

但这不是我在我的应用程序中想要的。我需要在计时器滴答发生时立即触发我所有的观察者更新方法。所以没有观察者必须等待。我希望这可以通过线程来完成。

所以,我将代码修改为,

// Fires the updates instantly
    public void Notify()
    {
      foreach (Observer o in _observers)
      {
        Threading.Thread oThread = new Threading.Thread(o.Update);
        oThread.Name = o.GetType().Name;
        oThread.Start();
      }
    }

但我心里有两个疑惑,

  1. 如果有 10 个观察者 我的定时器间隔是60秒 然后语句 new Thread() 将触发 600 次。

    是否高效并建议在每个计时器滴答时创建新线程?

  2. 如果我的观察者花费太多时间来完成他们的更新逻辑,即超过 60 秒怎么办。表示计时器滴答发生在观察者更新之前。我该如何控制?

我可以发布示例代码..如果需要...

我使用的代码..

using System;
using System.Collections.Generic;
using System.Timers;
using System.Text;
using Threading = System.Threading;
using System.ComponentModel;

namespace singletimers
{
  class Program
  {


    static void Main(string[] args)
    {
      DataPullerSubject.Instance.Attach(Observer1.Instance);
      DataPullerSubject.Instance.Attach(Observer2.Instance);
      Console.ReadKey();
    }
  }

  public sealed class DataPullerSubject
  {
    private static volatile DataPullerSubject instance;
    private static object syncRoot = new Object();
    public static DataPullerSubject Instance
    {
      get
      {
        if (instance == null)
        {
          lock (syncRoot)
          {
            if (instance == null)
              instance = new DataPullerSubject();
          }
        }

        return instance;
      }
    }

    int interval = 10 * 1000;
    Timer tmr;
    private List<Observer> _observers = new List<Observer>();

    DataPullerSubject()
    {
      tmr = new Timer();
      tmr.Interval = 1; // first time to call instantly
      tmr.Elapsed += new ElapsedEventHandler(tmr_Elapsed);
      tmr.Start();
    }

    public void Attach(Observer observer)
    {
      _observers.Add(observer);
    }

    public void Detach(Observer observer)
    {
      _observers.Remove(observer);
    }

    // Fires the updates instantly
    public void Notify()
    {
      foreach (Observer o in _observers)
      {
        Threading.Thread oThread = new Threading.Thread(o.Update);
        oThread.Name = o.GetType().Name;
        oThread.Start();
      }
    }

    private void tmr_Elapsed(object source, ElapsedEventArgs e)
    {
      tmr.Interval = interval;
      tmr.Stop(); // stop the timer until all notification triggered
      this.Notify();
      tmr.Start();//start again
    }
  }


  public abstract class Observer
  {
    string data;
    public abstract void Update();
    public virtual void GetDataFromDBAndSetToDataSet(string param)
    {
      Console.WriteLine("Processing for: " + param);
      data = param + new Random().Next(1, 2000);
      Threading.Thread.Sleep(10 * 1000);//long work
      Console.WriteLine("Data set for: " + param);
    }
  }


  public sealed class Observer1 : Observer
  {
    private static volatile Observer1 instance;
    private static object syncRoot = new Object();
    public static Observer1 Instance
    {
      get
      {
        if (instance == null)
        {
          lock (syncRoot)
          {
            if (instance == null)
              instance = new Observer1();
          }
        }

        return instance;
      }
    }
    Observer1()
    {
    }
    public override void Update()
    {
      base.GetDataFromDBAndSetToDataSet("Observer1");
    }

  }

  public sealed class Observer2 : Observer
  {
    private static volatile Observer2 instance;
    private static object syncRoot = new Object();
    public static Observer2 Instance
    {
      get
      {
        if (instance == null)
        {
          lock (syncRoot)
          {
            if (instance == null)
              instance = new Observer2();
          }
        }

        return instance;
      }
    }
    Observer2()
    {
    }
    public override void Update()
    {
      base.GetDataFromDBAndSetToDataSet("Observer2");
    }

  }
}

谢谢和亲切的问候。

【问题讨论】:

标签: c# .net multithreading oop design-patterns


【解决方案1】:
  • 不鼓励使用new Thread。使用Task 或Task&lt;T&gt;
  • 创建 Observable 模式框架的最佳尝试可能只接近于Rx。使用解决您提到的问题的方法(即,如果处理需要太多时间)。 Rx 将在定义可观察场景方面为您提供极大的灵活性。

【讨论】:

    【解决方案2】:

    1) 您可以通过 ThreadPool.QueueUserWorkItem 使用 ThreadPool 中的线程,也可以使用 Tasks

    2) 你必须synchronize your Methods。

    【讨论】:

      【解决方案3】:

      或者,观察者可以以非阻塞方式实现更新。 也就是说,Update 总是立即返回。然后,观察者对象负责在必要时在新线程中执行它们的工作。

      我不确定这对您的场景是否有帮助 - 我不知道您的“观察者”是什么,但也许您也不知道?

      【讨论】:

      • @Garen:我不明白你所说的“非阻塞方式”。在我的观察者中,我从数据库中获取数据。因此,这里可能需要几分钟以上的时间。这样做每个观察者都会更新他们的公共领域。
      • @thinkmmk - 我的意思是对 Update 的调用立即返回,并且观察者负责在单独的线程中启动 Update 操作。当更新已经在进行中时,观察者还负责排队或忽略更新调用。但这是否是一个好的决定取决于您的观察者 - 您希望避免在不同的“观察者”类中重复相同的逻辑。
      • 我已经用代码编辑了我的原始帖子。现在我正在为同步而苦苦挣扎,即不应该再次调用 obesever 'x',除非它完成了最后的工作。
      猜你喜欢
      • 2010-12-24
      • 1970-01-01
      • 2011-09-25
      • 2016-02-20
      • 2023-04-10
      • 2013-12-02
      • 2014-05-15
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多