【问题标题】:Batch process all items in ConcurrentBag批量处理 ConcurrentBag 中的所有项目
【发布时间】:2018-06-03 00:02:47
【问题描述】:

我有以下用例。多个线程正在创建收集在 ConcurrentBag 中的数据点。每 x 毫秒,一个消费者线程会查看自上次以来传入的数据点并对其进行处理(例如,对它们进行计数 + 计算平均值)。

以下代码或多或少代表了我想出的解决方案:

private static ConcurrentBag<long> _bag = new ConcurrentBag<long>();

static void Main()
{
    Task.Run(() => Consume());
    var producerTasks = Enumerable.Range(0, 8).Select(i => Task.Run(() => Produce()));
    Task.WaitAll(producerTasks.ToArray());
}

private static void Produce()
{
    for (int i = 0; i < 100000000; i++)
    {
        _bag.Add(i);
    }
}

private static void Consume()
{
    while (true)
    {
        var oldBag = _bag;
        _bag = new ConcurrentBag<long>();
        var average = oldBag.DefaultIfEmpty().Average();
        var count = oldBag.Count;
        Console.WriteLine($"Avg = {average}, Count = {count}");
        // Wait x ms
    }
}
  • ConcurrentBag 是适合这里工作的工具吗?
  • 切换包是实现清除新数据点列表然后处理旧数据点的正确方法吗?
  • 在 oldBag 上操作是否安全,或者当我迭代 oldBag 并且线程仍在添加项目时会遇到麻烦?
  • 我应该使用 Interlocked.Exchange() 来切换变量吗?

编辑

我猜上面的代码并不能很好地代表我想要实现的目标。所以这里有更多的代码来说明问题:

public class LogCollectorTarget : TargetWithLayout, ILogCollector
{
    private readonly List<string> _logMessageBuffer;

    public LogCollectorTarget()
    {
        _logMessageBuffer = new List<string>();
    }

    protected override void Write(LogEventInfo logEvent)
    {
        var logMessage = Layout.Render(logEvent);
        lock (_logMessageBuffer)
        {
            _logMessageBuffer.Add(logMessage);
        }
    }

    public string GetBuffer()
    {
        lock (_logMessageBuffer)
        {
            var messages =  string.Join(Environment.NewLine, _logMessageBuffer);
            _logMessageBuffer.Clear();
            return messages;
        }
    }
}

该类的目的是收集日志,以便将它们分批发送到服务器。每 x 秒调用一次 GetBuffer。这应该获取当前日志消息并清除缓冲区以获取新消息。它适用于锁,但由于它们非常昂贵,我不想锁定程序中的每个日志记录操作。所以这就是我想使用 ConcurrentBag 作为缓冲区的原因。但是当我调用 GetBuffer 时,我仍然需要切换或清除它,而不会丢失切换期间发生的任何日志消息。

【问题讨论】:

    标签: c# concurrency


    【解决方案1】:

    由于您只有一个消费者,因此您可以使用简单的 ConcurrentQueue 以自己的方式工作,而无需交换集合:

    public class LogCollectorTarget : TargetWithLayout, ILogCollector
    {
        private readonly ConcurrentQueue<string> _logMessageBuffer;
    
        public LogCollectorTarget()
        {
            _logMessageBuffer = new ConcurrentQueue<string>();
        }
    
        protected override void Write(LogEventInfo logEvent)
        {
            var logMessage = Layout.Render(logEvent);
            _logMessageBuffer.Enqueue(logMessage);
        }
    
        public string GetBuffer()
        {
            // How many messages should we dequeue?
            var count = _logMessageBuffer.Count;
    
            var messages = new StringBuilder();
    
            while (count > 0 && _logMessageBuffer.TryDequeue(out var message))
            {   
                messages.AppendLine(message);   
                count--;
            }       
    
            return messages.ToString();
        }
    }
    

    如果内存分配成为问题,您可以将它们出列到一个固定大小的数组并调用string.Join。这样,您就可以保证只进行两次分配(而如果初始缓冲区的大小不合适,StringBuilder 可以做更多的事情):

    public string GetBuffer()
    {
        // How many messages should we dequeue?
        var count = _logMessageBuffer.Count;
        var buffer = new string[count];
    
        for (int i = 0; i < count; i++)
        {   
            _logMessageBuffer.TryDequeue(out var message);
            buffer[i] = message;   
        }       
    
        return string.Join(Environment.NewLine, buffer);
    }
    

    【讨论】:

    • 谢谢,我希望有一个“DequeueAll”方法,但应该可以!
    【解决方案2】:

    ConcurrentBag 是适合这里工作的工具吗?

    它是适合工作的工具,这实际上取决于您要做什么以及为什么。您给出的示例非常简单,没有任何上下文,因此很难说。

    换包是实现清单的正确方法吗 新数据点然后处理旧数据点?

    答案是否定的,可能有很多原因。如果一个线程在你切换它时写入它会发生什么?

    在 oldBag 上操作是否安全? 遍历 oldBag 并且线程仍在添加项目?

    不,您只是复制了引用,这将无济于事。

    我应该使用 Interlocked.Exchange() 来切换变量吗?

    互锁方法是很棒的东西,但是这对你当前的问题没有帮助,它们是为了线程安全地访问整数类型值。您真的很困惑,您需要查找更多线程安全示例。


    但是,让我们为您指明正确的方向。忘记 ConcurrentBag 和那些花哨的课程。我的建议是从简单开始并使用锁定,以便您了解问题的本质。

    如果您希望多个任务/线程访问一个列表,您可以轻松使用lock 语句并保护对列表/数组的访问,这样其他讨厌的线程就不会修改它。

    显然您编写的代码是一个无意义的示例,我的意思是您只是将连续数字添加到列表中,并让另一个线程对它们进行平均。这根本不需要是消费者生产者,并且只是同步更有意义。

    在这一点上,我会向您指出可以让您实现这种模式的更好的架构,例如 Tpl Dataflow,但我担心这只是一种学习消费,不幸的是,您确实需要更多地阅读多线程并尝试更多示例在我们能够真正帮助您解决问题之前。

    【讨论】:

    • 感谢您的 cmets。上面的例子不能很好地说明问题。我编辑了问题
    猜你喜欢
    • 1970-01-01
    • 2020-10-09
    • 2011-07-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-06-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多