Последовательная передача. Реализация

Есть паттерн. Коротко: необходимость последовательной обработки групп сообщений, которые не блокируют друг друга.

Но не удалось найти понятный и простой пример реализации этого паттерна аля 'best practice'. По этому сделал решение на свой субективный взгляд отвечающее требованиям.

Хотелось бы узнать, сталкивался кто-либо с подобной проблемой? Как решили эту проблему? Ещё есть просьба проверить моё решение.

В моём случае стояла задача обрабатывать большой объём сообщений. Сообщения представляют из себя изменения какого-то события (EventId), в рамках одного события сообщения должны обрабатываться последовательно. События друг от друга не зависят.

Разделил ответственности по работе на Master/Slave. Master - занимается распределением задач и созданием/удалением Slave'ов. Slave - копит сообщения в очередь и обрабатывает их в том порядке, в котором получил.

    public class Master
    {
        private readonly IStrategy _strategy;
        private readonly TimeSpan _slaveLifetime;
        private readonly ConcurrentDictionary<int, Slave> _container;
        private readonly Timer _timer;

        public Master(
            IStrategy strategy,
            TimeSpan slaveLifetime,
            TimeSpan slaveClean)
        {
            _strategy = strategy;
            _slaveLifetime = slaveLifetime;
            _container = new ConcurrentDictionary<int, Slave>();
            _timer = new Timer(KillSlaves, null, slaveClean, slaveClean);
        }

        public void Send(Message message)
        {
            if (!_container.TryGetValue(message.EventId, out var slave))
            {
                slave = new Slave(_strategy);
            }

            slave.Timestamp = DateTime.UtcNow;
            slave.Send(message);

            _container.TryAdd(message.EventId, slave);
        }

        private void KillSlaves(object state)
        {
            var dif = DateTime.UtcNow - _slaveLifetime;

            foreach (var item in _container)
            {
                if (item.Value.Timestamp < dif)
                {
                    _container.TryRemove(item.Key, out _);
                }
            }
        }
    }


    public class Slave
    {
        private readonly IStrategy _strategy;
        private readonly ActionBlock<Message> _block;

        public Slave(IStrategy strategy)
        {
            _strategy = strategy;
            _block = new ActionBlock<Message>(Handle);
        }

        public DateTime Timestamp { get; set; }

        public void Send(Message message)
        {
            _block.Post(message);
        }

        private void Handle(Message message)
        {
            _strategy.Execute(message);
        }
    }

Ответы (0 шт):