Organizing execution of events in a distributed system and avoid deadlocks

Viewed 254

I have an issue prioritizing events in my system. I have a simple class which can subscribe to each others output

public interface INode<TIn, TOut> : IBaseNode
{
    event EventHandler<TOut> Output; 
    //Note: subscribe just calls node.Output += this.OnInput
    void Subscribe(IBaseNode node);
    void OnInput(object sender, TIn input)
} 

using this I can chain nodes together by subscribing to their output

CarDealerNode.Subscribe(NewModelNode);
LoggerNode.Subscribe(CarDealerNode);

My issue is that when an event fires off, it happens in an semi-undeterminstic breadth first manner. I would like to maintain the order of the execution of these events so I can prioritize event execution in a more dynamic manner.

My first impression is to use a some priority queue to sort the tasks, However this may cause issues because lower priority things may never execute

public class SynchronizationInfo
{
    public SyncPriority Priority { get; set; } = SyncPriority.Normal;
    public object Sender { get; set; }
    public DateTime Created { get; set; } = DateTime.Now;
    public Task Operation { get; set; }
}

public class SynchronizationContext
{
    public PriorityQueue<SynchronizationInfo> ExecutionQueue = new PriorityQueue<SynchronizationInfo>();
    //...
}

However I'm still having trouble of grasping a way to assure that dead locks won't occur, If something of a high priority is added at a quicker rate than the execution of that priority, lower priority events won't execute.

Additionally, just because something is over-lower priority doesn't mean everything of higher priority should go first, time is a big factor.

Is there a solid efficient recommended way of handing priority execution of tasks. In a way that no task experiences dead-lock, (e.g time increases priority in a way in which lower priorities are moved up to assure execution)?

3 Answers

Why have one queue when we can have more?

This works for any constant priority count although a low-enough number is preferred. (From what I can see you have an enum for them so there are likely just a few priorities). Also, we won't be using a priority queue but several ordinary queues.

  • Make a queue for every priority. The task gets registered to a queue according to it's priority. For every task store the creation timestamp as @Funk does.
  • When you wish to process next task, check the timestamp for the available element in every queue.
  • This allows you to detect long-overdue lower priority tasks and to increase their priority.

The increase of the priority can be done in several ways. For example:

  • Start the task execution directly when it is in the queue long enough. (E.g high_creation_time > medium_creation_time + C -> run the medium priority task)
  • Re-schedule the task to the queue with higher priority instead of running directly.

Which way is more suitable for you is a bit hard to tell.

The complexity of this approach:

  • Adding new task: O(1) - just add it to it's respective queue
  • Running a task: O(1) - assuming that we have constant number of priorities, this is just a matter of checking all the queues and finding the element that should run next.
  • Rescheduling a task(if applicable) - one push and one pop, thus O(1)

Is there a solid efficient recommended way of handing priority execution of tasks. In a way that no task experiences dead-lock, (e.g time increases priority in a way in which lower priorities are moved up to assure execution)?

The problem with the "time increases priority" approach is that the priority-queue needs to be recalculated all the time.

Let's review the normal use case for a priority-queue. The following represents a listing of the ordered items in the data structure:

  • { Priority = SyncPriority.High, Created = "2021-03-05 12:34:01", ... }
  • { Priority = SyncPriority.High, Created = "2021-03-05 12:34:04", ... }
  • { Priority = SyncPriority.High, Created = "2021-03-05 12:34:06", ... }
  • { Priority = SyncPriority.Normal, Created = "2021-03-05 12:34:02", ... }
  • { Priority = SyncPriority.Normal, Created = "2021-03-05 12:34:05", ... }
  • { Priority = SyncPriority.Low, Created = "2021-03-05 12:34:03", ... }

We can tell an event occurred every second, beginning with a High, Normal, Low,... priority. When the next High priority is added, we can insert it before the first item with Normal priority. The efficiency of the data structure is based upon the fact that the order of all the other items doesn't change.

If instead of the creation time, we add the passed time (ie a growing time interval) to the mix, the order would have to be recalculated on each request for the next highest priority item. The key basically becomes linear in time instead of constant.

To avoid this conundrum, and its inherent complexity, you could divide and conquer. Allowing one of a couple of simplifications, which employ a windowing system.

A max priority-queue exposes the following members:

  • Insert
  • RemoveMax

First simplification example, putting a bound on the number of elements in the priority-queue.

public class PriorityQueueWithMaxElements
{
    private readonly int _maxElementsInPriorityQueue;

    private readonly PriorityQueue<SynchronizationInfo> _priorityQueue;
    private readonly Queue<SynchronizationInfo> _backingQueue;

    public PriorityQueueWithMaxElements(int maxElementsInQueue)
    {
        _maxElementsInPriorityQueue = maxElementsInQueue;

        _priorityQueue = new PriorityQueue<SynchronizationInfo>();
        _backingQueue = new Queue<SynchronizationInfo>();
    }

    public void Insert(SynchronizationInfo info)
    {
        if (_backingQueue.Any() || _priorityQueue.Count == _maxElementsInPriorityQueue)
        {
            _backingQueue.Enqueue(info);
        }
        else
        {
            _priorityQueue.Insert(info);
        }
    }

    public SynchronizationInfo RemoveMax()
    {
        var max = _priorityQueue.RemoveMax();

        if (max == null && _backingQueue.Count > 0)
        {
            var numberOfItems = Math.Min(_maxElementsInPriorityQueue, _backingQueue.Count);
            for (var i = 0; i < numberOfItems; i++)
            {
                _priorityQueue.Insert(_backingQueue.Dequeue());
            }

            max = _priorityQueue.RemoveMax();
        }

        return max;
    }
}

Second example, putting a bound on the max time delay between the elements in the priority-queue.

public class PriorityQueueWithMaxDelay
{
    private readonly TimeSpan _maxDelay;
    
    private readonly PriorityQueue<SynchronizationInfo> _priorityQueue;
    private readonly Queue<SynchronizationInfo> _backingQueue;

    // DateTime of oldest item in the priority-queue.
    private DateTime? _baseDateTime;

    public PriorityQueueWithMaxDelay(TimeSpan maxDelay)
    {
        _maxDelay = maxDelay;

        _priorityQueue = new PriorityQueue<SynchronizationInfo>();
        _backingQueue = new Queue<SynchronizationInfo>();
    }

    public void Insert(SynchronizationInfo info)
    {
        if (_baseDateTime == null)
        {
            _baseDateTime = info.Created;
        }

        if (_backingQueue.Any() || !IsWithinDelay(info))
        {
            _backingQueue.Enqueue(info);
        }
        else
        {
            _priorityQueue.Insert(info);
        }
    }

    public SynchronizationInfo RemoveMax()
    {
        var max = _priorityQueue.RemoveMax();

        if (max == null && _backingQueue.Count > 0)
        {
            _baseDateTime = _backingQueue.Peek().Created;

            var numberOfItems = _backingQueue.TakeWhile(IsWithinDelay).Count();
            for (var i = 0; i < numberOfItems; i++)
            {
                _priorityQueue.Insert(_backingQueue.Dequeue());
            }

            max = _priorityQueue.RemoveMax();
        }

        if (_priorityQueue.Count == 0)
        {
            _baseDateTime = null;
        }


        return max;
    }

    public bool IsWithinDelay(SynchronizationInfo info)
        => info.Created - _baseDateTime < _maxDelay;
}

Note once a bound has been breached, the priority-queue is fully emptied before replenishing it again. This to avoid low priority items to remain in the queue indefinitely. I left multi-threading support as an exercise for the reader. Suffice to say that a performant implementation will be easier to achieve, then when you would have to recalculate the order of all items on each request.

Currently, the events are executed in the order they are registered. However, this is just the way it is implemented now, and I would not recommend anyone to rely on this behavior because the same may not be true in future versions.

.NET languages are hiding events implementation behind syntactical obstacles, although event handling is based on a relatively clearly defined delegate system.

Check the Delegate Class from here: Delegate Class

You can also change ordering even after attaching the events by detaching all handlers, and then re-attaching in desired order.

Note: You can override the default behaviour of events by changing the add and remove operations on your event to specify some other behaviour. You can then keep your event handlers in a list that you manage yourself and handle the firing order based on whatever rules you like.

Related