【问题标题】:C# queueing dependant tasks to be processed by a thread poolC# 将要由线程池处理的依赖任务排队
【发布时间】:2012-06-27 15:26:45
【问题描述】:

我想将需要按顺序(在每个流中)处理的多个流中的依赖任务排队。这些流可以并行处理。

具体来说,假设我需要两个队列,并且我希望每个队列中的任务按顺序处理。以下是用于说明所需行为的示例伪代码:

Queue1_WorkItem wi1a=...;

enqueue wi1a;

... time passes ...

Queue1_WorkItem wi1b=...;

enqueue wi1b; // This must be processed after processing of item wi1a is complete

... time passes ...

Queue2_WorkItem wi2a=...;

enqueue wi2a; // This can be processed concurrently with the wi1a/wi1b

... time passes ...

Queue1_WorkItem wi1c=...;

enqueue wi1c; // This must be processed after processing of item wi1b is complete

这是一个带有箭头的图表,说明了工作项之间的依赖关系:

问题是如何使用 C# 4.0/.NET 4.0 做到这一点?现在我有两个工作线程,每个队列一个,我为每个队列使用BlockingCollection<>。我想改为利用 .NET 线程池并让工作线程同时(跨流)处理项目,但在流中串行处理。换句话说,我希望能够指出例如 wi1b 取决于 wi1a 的完成,而不必跟踪完成并记住 wi1a,当 wi1b 到达时。换句话说,我只想说,“我想为 queue1 提交一个工作项,该工作项将与我已经为 queue1 提交的其他项串行处理,但可能与提交到其他队列的工作项并行处理”。

我希望这个描述是有道理的。如果没有,请随时在 cmets 中提问,我会相应地更新这个问题。

感谢阅读。

更新:

总结到目前为止“有缺陷”的解决方案,以下是我无法使用的答案部分的解决方案以及我无法使用它们的原因:

TPL 任务需要为 ContinueWith() 指定前面的任务。我不想在提交新任务时保留每个队列的先前任务的知识。

TDF ActionBlocks 看起来很有希望,但似乎发布到 ActionBlock 的项目是并行处理的。我需要对特定队列的项目进行连续处理。

更新 2:

RE:动作块

似乎将MaxDegreeOfParallelism 选项设置为1 会阻止并行处理提交给单个ActionBlock 的工作项。因此,似乎每个队列都有一个ActionBlock 解决了我的问题,唯一的缺点是这需要安装和部署 Microsoft 的 TDF 库,我希望有一个纯 .NET 4.0 解决方案。到目前为止,这是候选人接受的答案,除非有人能找到一种方法来使用纯 .NET 4.0 解决方案来做到这一点,该解决方案不会退化为每个队列的工作线程(我已经在使用)。

【问题讨论】:

  • 你看过Task/ContinueWith吗?
  • 我有并且我注意到 ContinueWith 需要了解先前的任务。我不想跟踪原始问题中指定的先前任务,部分原因是我必须按队列这样做。相反,我想从提交点“触发并忘记”并处理任务处理谓词中的错误和错误传播。换句话说,我想要提交时的最小状态——工作项和它应该提交到的队列。

标签: c# .net c#-4.0 .net-4.0


【解决方案1】:

我了解您有很多队列并且不想占用线程。每个队列可以有一个ActionBlock。 ActionBlock 可以自动完成您需要的大部分工作:它按顺序处理工作项,并且仅在工作未决时启动任务。当没有待处理的工作时,不会阻塞任何任务/线程。

【讨论】:

  • 不幸的是,根据您提供的链接中的文档,此解决方案需要尚未发布的 .NET 4.5。我正在寻找一种利用 .NET 4.0 的解决方案,如该问题所附标签中所指定的那样。
  • TPL 数据流现成可用 (msdn.microsoft.com/en-us/devlabs/gg585582.aspx)。
  • 我已经安装了 TPL Data flow 并在msdn.microsoft.com/en-us/library/hh462696(v=vs.110).aspx 尝试了示例代码。似乎发布到操作块的项目是并行处理的,而不是您的回复中提到的串行处理。这不是我想要的 - 每个队列的项目必须按顺序处理。
  • 我明白这一点。将 ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 1 } 传递给构造函数。我相信 ActionBlock 正是您想要的,并且会为您带来很少的工作。
  • 这似乎很有效,并且以安装 TDF 为代价完美地解决了我的问题,这是可以接受的。我会再玩一些,如果它有效,你会得到复选标记(除非有人给我一个纯 .NET 4.0 解决方案,满足我的问题的限制)。
【解决方案2】:

最好的方法是使用Task Parallel Library (TPL)Continuations。延续不仅允许您创建任务流,还可以处理您的异常。这是 TPL 的great introduction。但是给你一些想法......

您可以使用

启动 TPL 任务
Task task = Task.Factory.StartNew(() => 
{
    // Do some work here...
});

现在要在前面的任务完成(错误或成功)时启动第二个任务,您可以使用ContinueWith 方法

Task task1 = Task.Factory.StartNew(() => Console.WriteLine("Antecedant Task"));
Task task2 = task1.ContinueWith(antTask => Console.WriteLine("Continuation..."));

所以只要task1 完成、失败或被取消task2 'fires-up' 并开始运行。请注意,如果task1 在到达第二行代码之前完成,task2 将被安排立即执行。传递给第二个 lambda 的 antTask 参数是对前面任务的引用。有关更多详细示例,请参阅this link...

您还可以传递先前任务的延续结果

Task.Factory.StartNew<int>(() => 1)
    .ContinueWith(antTask => antTask.Result * 4)
    .ContinueWith(antTask => antTask.Result * 4)
    .ContinueWith(antTask =>Console.WriteLine(antTask.Result * 4)); // Prints 64.

注意。请务必阅读提供的第一个链接中的异常处理,因为这可能会使新手误入 TPL。

最后一件需要特别注意的是子任务。子任务是那些创建为AttachedToParent 的任务。在这种情况下,直到所有子任务都完成后,延续才会运行

TaskCreationOptions atp = TaskCreationOptions.AttachedToParent;
Task.Factory.StartNew(() =>
{
    Task.Factory.StartNew(() => { SomeMethod() }, atp);
    Task.Factory.StartNew(() => { SomeOtherMethod() }, atp); 
}).ContinueWith( cont => { Console.WriteLine("Finished!") });

我希望这会有所帮助。

编辑:你看过ConcurrentCollections,尤其是BlockngCollection&lt;T&gt;。因此,在您的情况下,您可能会使用类似

public class TaskQueue : IDisposable
{
    BlockingCollection<Action> taskX = new BlockingCollection<Action>();

    public TaskQueue(int taskCount)
    {
        // Create and start new Task for each consumer.
        for (int i = 0; i < taskCount; i++)
            Task.Factory.StartNew(Consumer);  
    }

    public void Dispose() { taskX.CompleteAdding(); }

    public void EnqueueTask (Action action) { taskX.Add(Action); }

    void Consumer()
    {
        // This seq. that we are enumerating will BLOCK when no elements
        // are avalible and will end when CompleteAdding is called.
        foreach (Action action in taskX.GetConsumingEnumerable())
            action(); // Perform your task.
    }
}

【讨论】:

  • 这里的缺陷是我不想在提交 task2 时保留 task1 的任何内存。我只想说 task2 应该在我提交的最后一个项目之后在 queue1 上提交,不管那是什么。因此,我既不知道也不想在提交下一个任务(task2)时跟踪前面的任务(task1)。
  • 看看我刚刚添加的子任务部分。这可以帮助您实现您所需要的。万事如意...
  • SomeMethod()SomeOtherMethod() 是串行执行还是并行执行?根据我的需要,它们必须按顺序执行。此外,我没有同时提供所有任务。工作项目进来了,我需要安排它们,所以我不能用多个项目做一个StartNew()
  • 它们是“并行”执行的(尽管不是严格意义上的),在这种情况下它们是并发的,因为没有task.Wait() 调用。要按顺序执行这些操作,请添加Task taskA = Task.Factiory.StartNew...;,然后添加taskA.Wait();。我添加了一些关于ConcurrentCollections() 的内容。我不是专家,但我已经使用这些在一个小型应用程序中进行了一个小型排队过程,并且效果很好。这可能会对您有所帮助 - 或者至少为您提供另一种说服方式。万事如意...
  • 不幸的是,taskA.Wait() 为我扼杀了交易。另外,正如原始问题中提到的,我目前已经在使用工作线程和来自并发集合的BlockingCollection。您的 TaskQueue 解决方案看起来很有希望,但实际上您正在创建长时间运行的任务,这些任务在很长一段时间内占用线程池中的线程 - 即,此方法等效于手动工作线程。如果我错了,请纠正我。
【解决方案3】:

基于 TPL 的 .NET 4.0 解决方案是可能的,同时隐藏了它需要将父任务存储在某处的事实。例如:

class QueuePool
{
    private readonly Task[] _queues;

    public QueuePool(int queueCount)
    { _queues = new Task[queueCount]; }

    public void Enqueue(int queueIndex, Action action)
    {
        lock (_queues)
        {
           var parent = _queue[queueIndex];
           if (parent == null)
               _queues[queueIndex] = Task.Factory.StartNew(action);
           else
               _queues[queueIndex] = parent.ContinueWith(_ => action());
        }
    }
}

这是对所有队列使用单个锁来说明这个想法。然而,在生产代码中,我会为每个队列使用一个锁来减少争用。

【讨论】:

  • 这不是一个好的解决方案,因为如果父任务已经完成,ContinueWith 将启动锁内的操作。
【解决方案4】:

看起来您已经拥有的设计很好并且有效。您的工作线程(每个队列一个)是长时间运行的,因此如果您想改用 Task 的,请指定 TaskCreationOptions.LongRunning 以便获得一个专用的工作线程。

但这里实际上不需要使用 ThreadPool。它不会为长期工作提供很多好处。

【讨论】:

  • 线程运行时间长,但任务本身很短。此外,如果我将其扩展到许多队列,它会崩溃,因为我可能会超额订阅,因为所有队列都是同时处理的。使用线程池,我不必担心这一点,因为并行队列中的任务可以由线程池并行处理,而上下文切换最少。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2012-02-12
  • 1970-01-01
  • 2016-04-20
  • 2022-01-18
  • 2011-07-18
  • 2010-12-18
  • 2015-03-27
相关资源
最近更新 更多