【问题标题】:Process a ConcurrentStack when not empty?不为空时处理 ConcurrentStack?
【发布时间】:2015-05-29 14:01:37
【问题描述】:

我有一个要向其中倾倒项目的 ConcurrentStack。当堆栈不为空时,一次处理这些项目的好方法是什么?我想以一种在未处理堆栈时不会占用 CPU 周期的方式来执行此操作。

我目前得到的基本上是这个,它似乎不是一个理想的解决方案。

private void AddToStack(MyObj obj)
{
    stack.Push(obj);
    HandleStack();
}

private void HandleStack()
{
    if (handling)
        return;

    Task.Run( () =>
    {
        lock (lockObj)
        {
            handling = true;
            if (stack.Any())
            {
                //handle whatever is on top of the stack
            }
            handling = false;
        }
    }
}

所以 bool 存在,因此多个线程不会在等待锁时得到备份。但我不希望同时处理堆栈的多个事物因此锁定。因此,如果两个单独的线程确实最终同时调用 HandleStack 并通过了布尔值,那么锁就在那里,所以两个线程都不会同时遍历堆栈。但是一旦第二个通过锁,堆栈将是空的并且不做任何事情。所以这最终给了我我想要的行为。

所以实际上我只是在 ConcurrentStack 周围编写一个伪并发包装器,而且似乎必须有一种不同的方法来实现这一点。想法?

【问题讨论】:

  • 您真的需要按顺序处理每个元素吗?如果它们可以同时处理,只需委托给线程池,该线程池将有效地处理其工作队列。
  • 你不需要锁。它是一个 ConcurrentStack,它可以被多个线程修改。如果您真的想在等待时阻塞,请使用 BlockingCollection。默认情况下它使用 ConcurrentQueue 但您可以指定不同的并发集合,例如 ConcurrentStack
  • @PanagiotisKanavos 我知道它是由多个线程修改的。我希望添加多个线程,但只能从中获取一个。这就是为什么我要锁定 pop(在“//handle whatever...”后面混淆)而不是 push。
  • @BenManes 是的,它们特别不能同时处理。
  • @claudekennilol 你想做什么?如果您只想要 一个 消费者,请不要添加多个 - 例如使用 ActionBlock 或消费者的单例实例。否则,您将获得循环处理,所有消费者都等到堆栈清空

标签: c# multithreading concurrency


【解决方案1】:

ConcurrentStack<T> 是实现IProducerConsumerCollection<T> 的集合之一,因此可以被BlockingCollection<T> 包装。 BlockingCollection<T> 有几个方便成员用于常见操作,例如“在堆栈不为空时使用”。例如,您可以循环调用TryTake。或者,您可以使用GetConsumingEnumerable

private BlockingCollection<MyObj> stack;
private Task consumer;

Constructor()
{
  stack = new BlockingCollection<MyObj>(new ConcurrentStack<MyObj>());
  consumer = Task.Run(() =>
  {
    foreach (var myObj in stack.GetConsumingEnumerable())
    {
      ...
    }
  });
}

private void AddToStack(MyObj obj)
{
  stack.Add(obj);
}

【讨论】:

  • 谢谢,这绝对是我需要的。我对消费者任务的工作方式感到困惑。我在循环中选择了 TryPop,因为我不想要恒定的 cpu 周期。 GetConsumingEnumerable 与此有何不同(即堆栈为空时发生了什么/添加某些内容时触发它的原因)?
  • 只要底层集合为空,TryPopMoveNext(在消费枚举上调用时)都会阻塞调用线程。我不知道确切的实现,但它在逻辑上是一个监视器。所以被阻塞的线程处于等待状态,不消耗CPU。
【解决方案2】:

您可以考虑使用Microsoft TPL Dataflow 来做这种事情。

这是一个简单的示例,展示了如何创建队列。试一试,调整MaxDegreeOfParallelismBoundedCapacity 的设置,看看会发生什么。

对于您的示例,如果您不希望多个线程同时处理一个数据项,我认为您需要将 MaxDegreeOfParallelism 设置为 1。

(注意:您需要使用 .Net 4.5x 并使用 Nuget 为项目安装 TPL Dataflow。)

还可以阅读Stephen Cleary's blog about TPL

using System;
using System.Threading;
using System.Threading.Tasks;
using System.Threading.Tasks.Dataflow;

namespace SimpleTPL
{
    class MyObj
    {
        public MyObj(string data)
        {
            Data = data;
        }

        public readonly string Data;
    }

    class Program
    {
        static void Main()
        {
            var queue = new ActionBlock<MyObj>(data => process(data), actionBlockOptions());
            var task = queueData(queue);

            Console.WriteLine("Waiting for task to complete.");
            task.Wait();
            Console.WriteLine("Completed.");
        }

        private static void process(MyObj data)
        {
            Console.WriteLine("Processing data " + data.Data);
            Thread.Sleep(200); // Simulate load.
        }

        private static async Task queueData(ActionBlock<MyObj> executor)
        {
            for (int i = 0; i < 20; ++i)
            {
                Console.WriteLine("Queuing data " + i);
                MyObj data = new MyObj(i.ToString());

                await executor.SendAsync(data);
            }

            Console.WriteLine("Indicating that no more data will be queued.");

            executor.Complete(); // Indicate that no more items will be queued.

            Console.WriteLine("Waiting for queue to empty.");

            await executor.Completion; // Wait for executor queue to empty.
        }

        private static ExecutionDataflowBlockOptions actionBlockOptions()
        {
            return new ExecutionDataflowBlockOptions
            {
                MaxDegreeOfParallelism = 4,
                BoundedCapacity        = 8
            };
        }
    }
}

【讨论】:

    【解决方案3】:

    看起来你想要一个典型的生产者消费者。

    我建议使用自动重置事件

    当堆栈为空时让您的消费者等待。调用生产者方法时调用Set。

    阅读本帖

    Fast and Best Producer/consumer queue technique BlockingCollection vs concurrent Queue

    【讨论】:

    • BlockingCollection 不是已经这样做了吗? BlockingCollection 上的阻塞与 AutoResetEvent 上的阻塞有什么区别?
    • @PanagiotisKanavos 我实际上并没有建议一起使用它们,但我可以看到你从哪里得到这种印象。我的回答有点笨拙。
    • 你误会了。既然可以只使用 BlockingCollection,为什么还要将 AutoResetEvent 与 ConcurrentStack 一起使用?
    • 哦,对了……当然。完全一致。但是“当你可以使用时为什么要使用”的论点可能会有点循环。这就是发展的美/本质。做某事总是有更好的方式或首选方式,但当涉及到主观性时,它们并不总是相同的。
    • 这不是主观问题 - BlockingCollection uses a semaphore itself 等待发布者。仅使用更多代码执行与 BCL 库相同的操作并不是最优的。最佳是例如。使用 Interlocked.CompareExchange 检查和设置标志而不阻塞
    猜你喜欢
    • 1970-01-01
    • 2011-12-19
    • 1970-01-01
    • 2018-10-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多