【问题标题】:Managing high/low priority threads in .net在 .net 中管理高/低优先级线程
【发布时间】:2013-02-11 21:15:49
【问题描述】:

场景如下: 有几个低优先级线程可以被高优先级线程中断。每当高优先级线程要求低优先级线程暂停时,它们将进入Wait 状态(如果它们尚未处于等待状态)。然而,当一个高优先级线程发出低优先级线程可以Resume 的信号时,低优先级线程不应恢复,直到所有请求低优先级线程暂停的高优先级线程都同意。

为了解决这个问题,我在计数器变量中跟踪从高优先级线程到低优先级线程的Pause() 调用。每当高优先级线程向低优先级线程请求Pause()时,计数器的值就加1。如果加1后计数器的值为1,则表示该线程不在Wait中,所以要求它进入Wait 状态。否则只需增加 counter 值。相反,当高优先级线程调用Resume()时,我们将counter的值递减,如果递减后的值为0,则表示低优先级线程现在可以Resume。

这是我的问题的简化实现。带有Interlocked.XXX的if语句内部的比较操作不正确,即

if (Interlocked.Increment(ref _remain) == 1)

,因为读取/修改和比较操作不是原子的。

我在这里缺少什么?我不想使用线程优先级。

using System;
using System.Collections.Generic;
using System.Threading;

namespace TestConcurrency
{

// I borrowed this class from Joe Duffy's blog and modified it
public class LatchCounter
{
 private long _remain;
 private EventWaitHandle m_event;
 private readonly object _lockObject;

public LatchCounter()
{
    _remain = 0;
    m_event = new ManualResetEvent(true);
    _lockObject = new object();
}

public void Check()
{
    if (Interlocked.Read(ref _remain) > 0)
    {
        m_event.WaitOne();
    }
}

public void Increment()
{
    lock(_lockObject)
    {
       if (Interlocked.Increment(ref _remain) == 1)
           m_event.Reset();
    }
}

public void Decrement()
{
    lock(_lockObject)
    {
       // The last thread to signal also sets the event.
       if (Interlocked.Decrement(ref _remain) == 0)
           m_event.Set();
    }
}
}



public class LowPriorityThreads
{
private List<Thread> _threads;
private LatchCounter _latch;
private int _threadCount = 1;

internal LowPriorityThreads(int threadCount)
{
    _threadCount = threadCount;
    _threads = new List<Thread>();
    for (int i = 0; i < _threadCount; i++)
    {
        _threads.Add(new Thread(ThreadProc));
    }

    _latch = new CountdownLatch();
}


public void Start()
{
    foreach (Thread t in _threads)
    {
        t.Start();
    }
}

void ThreadProc()
{
    while (true)
    {
        //Do something
        Thread.Sleep(Rand.Next());
        _latch.Check();
    }
}

internal void Pause()
{
    _latch.Increment();
}

internal void Resume()
{
    _latch.Decrement();
}
}


public class HighPriorityThreads
{
private Thread _thread;
private LowPriorityThreads _lowPriorityThreads;

internal HighPriorityThreads(LowPriorityThreads lowPriorityThreads)
{
    _lowPriorityThreads = lowPriorityThreads;
    _thread = new Thread(RandomlyInterruptLowPriortyThreads);
}

public void Start()
{
    _thread.Start();
}

void RandomlyInterruptLowPriortyThreads()
{
    while (true)
    {
        Thread.Sleep(Rand.Next());

        _lowPriorityThreads.Pause();

        Thread.Sleep(Rand.Next());
        _lowPriorityThreads.Resume();
    }
}
}

 class Program
 {
  static void Main(string[] args)
  {
    LowPriorityThreads lowPriorityThreads = new LowPriorityThreads(3);
    HighPriorityThreads highPriorityThreadOne = new HighPriorityThreads(lowPriorityThreads);
    HighPriorityThreads highPriorityThreadTwo = new HighPriorityThreads(lowPriorityThreads);

    lowPriorityThreads.Start();
    highPriorityThreadOne.Start();
    highPriorityThreadTwo.Start();
}
}


class Rand
{
internal static int Next()
{
    // Guid idea has been borrowed from somewhere on StackOverFlow coz I like it
    return new System.Random(Guid.NewGuid().GetHashCode()).Next() % 30000;
}
}

【问题讨论】:

  • 你为什么不能直接修改 Check 来做m_event.WaitOne() 而不做其他事情?
  • 不使用 Thread.Priority 是一个严重的错误。在奶牛回家之前,您将调试死锁。
  • 这个要求有点“不对劲”,但我现在还不能完全确定。
  • @usr:我可能已经等不及了。如果不止一个高优先级线程处于运行状态,那么我必须等待所有高优先级线程发出信号
  • “我可能不会等待”:在这种情况下,将始终设置事件,因此不会发生等待。这对我来说是正确的。

标签: c# multithreading concurrency interlocked-increment


【解决方案1】:

我不知道您的要求,因此我不会在这里讨论它们。 就实现而言,我将引入一个“调度程序”类,该类将处理线程间交互并充当“可运行”对象的工厂。

实施当然是非常粗略的,而且很容易受到批评。

class Program
{
    static void Main(string[] args)
    {
        ThreadDispatcher td=new ThreadDispatcher();
        Runner r1 = td.CreateHpThread(d=>OnHpThreadRun(d,1));
        Runner r2 = td.CreateHpThread(d => OnHpThreadRun(d, 2));

        Runner l1 = td.CreateLpThread(d => Console.WriteLine("Running low priority thread 1"));
        Runner l2 = td.CreateLpThread(d => Console.WriteLine("Running low priority thread 2"));
        Runner l3 = td.CreateLpThread(d => Console.WriteLine("Running low priority thread 3"));


        l1.Start();
        l2.Start();
        l3.Start();

        r1.Start();
        r2.Start();

        Console.ReadLine();

        l1.Stop();
        l2.Stop();
        l3.Stop();

        r1.Stop();
        r2.Stop();
    }

    private static void OnHpThreadRun(ThreadDispatcher d,int number)
    {
        Random r=new Random();
        Thread.Sleep(r.Next(100,2000));
        d.CheckedIn();
        Console.WriteLine(string.Format("*** Starting High Priority Thread {0} ***",number));
        Thread.Sleep(r.Next(100, 2000));
        Console.WriteLine(string.Format("+++ Finishing High Priority Thread {0} +++", number));
        Thread.Sleep(300);
        d.CheckedOut();           
    }
}

public abstract class Runner
{
    private Thread _thread;
    protected readonly Action<ThreadDispatcher> _action;
    private readonly ThreadDispatcher _dispathcer;
    private long _running;
    readonly ManualResetEvent _stopEvent=new ManualResetEvent(false);
    protected Runner(Action<ThreadDispatcher> action,ThreadDispatcher dispathcer)
    {
        _action = action;
        _dispathcer = dispathcer;
    }

    public void Start()
    {
        _thread = new Thread(OnThreadStart);
        _running = 1;
        _thread.Start();
    }

    public void Stop()
    {
        _stopEvent.Reset();
        Interlocked.Exchange(ref _running, 0);
        _stopEvent.WaitOne(2000);
        _thread = null;
        Console.WriteLine("The thread has been stopped.");

    }
    protected virtual void OnThreadStart()
    {
        while (Interlocked.Read(ref _running)!=0)
        {
            OnStartWork();
            _action.Invoke(_dispathcer);
            OnFinishWork();
        }
        OnFinishWork();
        _stopEvent.Set();
    }

    protected abstract void OnStartWork();
    protected abstract void OnFinishWork();
}

public class ThreadDispatcher
{
    private readonly ManualResetEvent _signal=new ManualResetEvent(true);
    private int _hpCheckedInThreads;
    private readonly object _lockObject = new object();

    public void CheckedIn()
    {
        lock(_lockObject)
        {
            _hpCheckedInThreads++;
            _signal.Reset();
        }
    }
    public void CheckedOut()
    {
        lock(_lockObject)
        {
            if(_hpCheckedInThreads>0)
                _hpCheckedInThreads--;
            if (_hpCheckedInThreads == 0)
                _signal.Set();
        }
    }

    private class HighPriorityThread:Runner 
    {
        public HighPriorityThread(Action<ThreadDispatcher> action, ThreadDispatcher dispatcher) : base(action,dispatcher)
        {
        }

        protected override void OnStartWork()
        {
        }

        protected override void OnFinishWork()
        {
        }
    }
    private class LowPriorityRunner:Runner
    {
        private readonly ThreadDispatcher _dispatcher;
        public LowPriorityRunner(Action<ThreadDispatcher> action, ThreadDispatcher dispatcher)
            : base(action, dispatcher)
        {
            _dispatcher = dispatcher;
        }

        protected override void OnStartWork()
        {
            Console.WriteLine("LP Thread is waiting for a signal.");
            _dispatcher._signal.WaitOne();
            Console.WriteLine("LP Thread got the signal.");
        }

        protected override void OnFinishWork()
        {

        }
    }

    public Runner CreateLpThread(Action<ThreadDispatcher> action)
    {
        return new LowPriorityRunner(action, this);
    }

    public Runner CreateHpThread(Action<ThreadDispatcher> action)
    {
        return new HighPriorityThread(action, this);
    }
}

}

【讨论】:

    猜你喜欢
    • 2014-11-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-04-19
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多