Последовательная передача. Реализация
Есть паттерн. Коротко: необходимость последовательной обработки групп сообщений, которые не блокируют друг друга.
Но не удалось найти понятный и простой пример реализации этого паттерна аля '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);
}
}