【发布时间】: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