生产者/消费者模式与批量生产者

Ant*_*ift 4 c# multithreading task thread-safety task-parallel-library

我正在尝试使用多个生产者和一个消费者来实现一个相当简单的生产者/消费者风格的应用程序.

研究让我进入了BlockingCollection<T>有用的领域,并允许我实现一个长期运行的消费者任务,如下所示:

var c1 = Task.Factory.StartNew(() =>
{
    var buffer = new List<int>(BATCH_BUFFER_SIZE);

    foreach (var value in blockingCollection.GetConsumingEnumerable())
    {
        buffer.Add(value);
        if (buffer.Count == BATCH_BUFFER_SIZE)
        {
            ProcessItems(buffer);
            buffer.Clear();
        }
    }
});
Run Code Online (Sandbox Code Playgroud)

该ProcessItems函数将缓冲区提交给数据库,并且它可以批量工作.然而,该解决方案是次优的.在男爵生产期间,可能需要一段时间才能填充缓冲区,这意味着数据库已过期.

更理想的解决方案是在30秒计时器上运行任务或foreach在超时时短路.

我跑了计时器的想法,想出了这个:

syncTimer = new Timer(new TimerCallback(TimerElapsed), blockingCollection, 5000, 5000);

private static void TimerElapsed(object state)
{
    var buffer = new List<int>();
    var collection = ((BlockingCollection<int>)state).GetConsumingEnumerable();

    foreach (var value in collection)
    {
        buffer.Add(value);
    }

    ProcessItems(buffer);
    buffer.Clear();
}
Run Code Online (Sandbox Code Playgroud)

这有明显的问题,foreach将被阻止直到结束,打败计时器的目的.

任何人都可以提供指导吗?我基本上需要BlockingCollection定期快照并处理内容清除它.也许BlockingCollection是错误的类型?

Ani*_*Ani 6

而不是GetConsumingEnumerable在计时器回调中使用,而是使用其中一种方法,将结果添加到列表中,直到它返回false或者您已达到满意的批量大小.

BlockingCollection.TryTake方法(T) - 可能你需要的,你根本不想进一步等待.

BlockingCollection.TryTake方法(T,Int32)

BlockingCollection.TryTake方法(T,TimeSpan)

您可以轻松地将其提取到扩展中(未经测试):

public static IList<T> Flush<T>
(this BlockingCollection<T> collection, int maxSize = int.MaxValue)
{
     // Argument checking.

     T next;
     var result = new List<T>();

     while(result.Count < maxSize && collection.TryTake(out next))
     {
         result.Add(next);
     }

     return result;
}
Run Code Online (Sandbox Code Playgroud)

  • @Andrew:这取决于所使用的Timer实现,该实现不在此答案的范围之内,该实现仅与阻塞集合本身有关。 (2认同)