How do I implement this "WorkerChain" functionality in .NET?

EDIT . It occurred to me too late (?) That all the code I posted in my first update to this question was too big for most readers. I really went ahead and wrote a blog post about this topic for anyone else who wants to read it.

In the meantime, I left the original question in place to briefly talk about the problem I would like to solve.

I also want to point out that the code I posted (on my blog) has worked great in testing so far. But I'm still interested in any feedback people are willing to give me on how clean / timid / fulfilled. *

* I love how this word doesn't really mean what we think , but developers use it all the time anyway.


Original question

Sorry for the vague title of the question - not sure how to encapsulate what I'm asking below, succinctly. (If someone with edit rights might think of a more descriptive title, feel free to change it.)

I need this behavior. I represent a working class that accepts one delegate task in its constructor (for simplicity I would make it immutable - no more tasks were added after instantiation). I'll call this assignment T

. The class should have a simple method, for example GetToWork

, that will demonstrate this behavior:

  • If the worker doesn't start T

    , then he will start doing it right now.
  • If the worker is currently running T

    , then as soon as it is finished it will start again T

    .
  • GetToWork

    can be called any number of times when the worker is working T

    ; a simple rule of thumb is that during any execution T

    , if GetToWork

    called at least once, it T

    will run again upon completion (and then if GetToWork

    called while T

    running at that time, it will repeat again, etc.).

It's pretty easy now with a boolean switch. But this class should be thread safe , by which I mean that steps 1 and 2 above should contain atomic operations (at least I think they do).

There is an additional layer of complexity. I need a "work chain" class that will be made up of many of these workers linked to each other. As soon as the first worker finishes the job, he essentially calls the GetToWork

worker after him; Meanwhile, if his own GetToWork

was called, it will restart as well. Logically, the call GetToWork

in the chain is essentially the same as the call GetToWork

for the first worker in the chain (I would fully assume that worker networks are not public).

One way to imagine what this hypothetical "work chain" would look like is to compare it to a team in a relay. Suppose there are four runners W1

through W4

, and let the chain be called C

. If I call C.StartWork()

, the following should happen:

  • If it W1

    is at the starting point (i.e. does nothing), it will start working towards W2

    .
  • If it is W1

    already working towards W2

    (i.e. completing its task), then as soon as it reaches W2

    , it will signal W2

    to start, immediately return to the starting point and, since Calling StartWork

    , start working towards direction again W2

    .
  • When it W1

    reaches the starting point W2

    , it will immediately return to its starting point.
    • If he W2

      just sits, he will immediately start working in the direction W3

      .
    • If it W2

      is not already working in the direction W3

      , then it W2

      will just walk again as soon as it reaches W3

      and returns to the starting point.

The above is probably a little confusing and poorly written. But hopefully you get the basic idea. Obviously, these workers will work in their own streams.

Also, I think maybe this feature already exists somewhere? If this is the case, definitely let me know!

+2


a source to share


3 answers


Use semaphores. Each worker is a thread with the following code (pseudocode):

WHILE(TRUE)
    WAIT_FOR_SEMAPHORE(WORKER_ID) //The semaphore for the current worker
    RESET_SEMAPHORE(WORKER_ID)
    /* DO WORK */
    POST_SEMAPHORE(NEXT_WORKER_ID) //The semaphore for the next worker
END

      



With a non-zero semaphore, it means that someone has signaled the current thread to do work. After a non-zero semaphore is received in its input line, it resets the semaphore (marked as not signaled by anyone), do the work (meanwhile the semaphore can be sent again), and dispatch the semaphore for the next worker. This story repeats itself in the next worker (s).

+1


a source


A naive implementation that you can get some mileage with.

Note:

I understand that scalar types i.e. bool flags that control execution are atomic assignment making them thread safe as you need / need in this scenario.



There are much more complex possibilities associated with semaphores and other strategies, but if simple works ...

using System;
using System.Threading;

namespace FlaggedWorkerChain
{
    internal class Program
    {
        private static void Main(string[] args)
        {
            FlaggedChainedWorker innerWorker = new FlaggedChainedWorker("inner", () => Thread.Sleep(1000), null);
            FlaggedChainedWorker outerWorker = new FlaggedChainedWorker("outer", () => Thread.Sleep(500), innerWorker);

            Thread t = new Thread(outerWorker.GetToWork);
            t.Start();

            // flag outer to do work again
            outerWorker.GetToWork();

            Console.WriteLine("press the any key");
            Console.ReadKey();
        }
    }

    public sealed class FlaggedChainedWorker
    {
        private readonly string _id;
        private readonly FlaggedChainedWorker _innerWorker;
        private readonly Action _work;
        private bool _busy;
        private bool _flagged;

        public FlaggedChainedWorker(string id, Action work, FlaggedChainedWorker innerWorker)
        {
            _id = id;
            _work = work;
            _innerWorker = innerWorker;
        }

        public void GetToWork()
        {
            if (_busy)
            {
                _flagged = true;
                return;
            }

            do
            {
                _flagged = false;
                _busy = true;
                Console.WriteLine(String.Format("{0} begin", _id));

                _work.Invoke();

                if (_innerWorker != null)
                {
                    _innerWorker.GetToWork();
                }
                Console.WriteLine(String.Format("{0} end", _id));

                _busy = false;
            } while (_flagged);
        }
    }
}

      

+1


a source


It seems to me that you are too insulting. I've written these "pipelined" classes before; all you need is a queue of workers, each with a wait descriptor that is signaled when the action completes.

public class Pipeline : IDisposable
{
    private readonly IEnumerable<Stage> stages;

    public Pipeline(IEnumerable<Action> actions)
    {
        if (actions == null)
            throw new ArgumentNullException("actions");
        stages = actions.Select(a => new Stage(a)).ToList();
    }

    public Pipeline(params Action[] actions)
        : this(actions as IEnumerable<Action>)
    {
    }

    public void Dispose()
    {
        foreach (Stage stage in stages)
            stage.Dispose();
    }

    public void Start()
    {
        foreach (Stage currentStage in stages)
            currentStage.Execute();
    }

    class Stage : IDisposable
    {
        private readonly Action action;
        private readonly EventWaitHandle readyEvent;

        public Stage(Action action)
        {
            this.action = action;
            this.readyEvent = new AutoResetEvent(true);
        }

        public void Dispose()
        {
            readyEvent.Close();
        }

        public void Execute()
        {
            readyEvent.WaitOne();
            action();
            readyEvent.Set();
        }
    }
}

      

And here's a test program that you can use to verify that the actions are always executed in the correct order and only one of the actions can be performed at once:

class Program
{
    static void Main(string[] args)
    {
        Action firstAction = GetTestAction(1);
        Action secondAction = GetTestAction(2);
        Action thirdAction = GetTestAction(3);
        Pipeline pipeline = new Pipeline(firstAction, secondAction, thirdAction);
        for (int i = 0; i < 10; i++)
        {
            ThreadPool.QueueUserWorkItem(s => pipeline.Start());
        }
    }

    static Action GetTestAction(int index)
    {
        return () =>
        {
            Console.WriteLine("Action started: {0}", index);
            Thread.Sleep(100);
            Console.WriteLine("Action finished: {0}", index);
        };
    }
}

      

Short, simple, completely thread safe.

If for some reason you need to start from a specific step in the chain, then you can simply add an overload for Start

:

public void Start(int index)
{
    foreach (Stage currentStage in stages.Skip(index + 1))
        currentStage.Execute();
}

      


Edit

Based on the comments, I think a few minor changes to the inner class Stage

should be enough to get the behavior you want. We just need to add the "queued" event in addition to the "ready" event.

    class Stage : IDisposable
    {
        private readonly Action action;
        private readonly EventWaitHandle readyEvent;
        private readonly EventWaitHandle queuedEvent;

        public Stage(Action action)
        {
            this.action = action;
            this.readyEvent = new AutoResetEvent(true);
            this.queuedEvent = new AutoResetEvent(true);
        }

        public void Dispose()
        {
            readyEvent.Close();
        }

        private bool CanExecute()
        {
            if (readyEvent.WaitOne(0, true))
                return true;
            if (!queuedEvent.WaitOne(0, true))
                return false;
            readyEvent.WaitOne();
            queuedEvent.Set();
            return true;
        }

        public bool Execute()
        {
            if (!CanExecute())
                return false;
            action();
            readyEvent.Set();
            return true;
        }
    }

      

Also change the pipeline method Start

to break if the stage cannot complete (i.e. already enqueued):

public void Start(int index)
{
    foreach (Stage currentStage in stages.Skip(index + 1))
        if (!currentStage.Execute())
            break;
}

      

The concept here is pretty simple, once again:

  • The stage first tries to get the ready state immediately. If it succeeds, it starts up.
  • If it cannot get the ready state (that is, the task has already started), it tries to get the state of the queue.
    • If it receives a state in the queue, it expects the ready state to become available and then releases the state in the queue.
    • If he cannot get the state in the queue, then he will refuse.

I read your question and comments again and I am sure this is exactly what you are trying to do and gives the best compromise between security, bandwidth and throttling.

Since ThreadPool

it can sometimes take a while for a response, you should delay the delay in the test program 1000

instead 100

if you really want to see the "gaps".

+1


a source







All Articles