【问题标题】:Ensure a long running task is only fired once and subsequent request are queued but with only one entry in the queue确保长时间运行的任务只触发一次,后续请求排队但队列中只有一个条目
【发布时间】:2014-08-07 19:13:24
【问题描述】:

我有一个计算密集型方法Calculate,它可能会运行几秒钟,请求来自多个线程。

只应执行一个Calculate,后续请求应排队等待初始请求完成。如果已经有一个请求排队,那么可以丢弃后续请求(因为排队的请求就足够了)

似乎有很多潜在的解决方案,但我只需要最简单的。

更新:这是我的初步尝试:

private int _queueStatus;
private readonly object _queueStatusSync = new Object();

public void Calculate()
{
    lock(_queueStatusSync)
    {
        if(_queueStatus == 2) return;
        _queueStatus++;
        if(_queueStatus == 2) return;
    }
    for(;;)
    {
        CalculateImpl();
        lock(_queueStatusSync)
            if(--_queueStatus == 0) return;

    }
}

private void CalculateImpl()
{
    // long running process will take a few seconds...
}

【问题讨论】:

    标签: c# .net multithreading asynchronous async-await


    【解决方案1】:

    IMO 最简单、最干净的解决方案是使用 TPL Dataflow(一如既往)和 BufferBlock 作为队列。 BufferBlock 是线程安全的,支持async-await,更重要的是,TryReceiveAll 可以一次获取所有项目。它还具有OutputAvailableAsync,因此您可以异步等待将项目发布到缓冲区。当发布多个请求时,您只需选择最后一个,然后忘记其余的:

    var buffer = new BufferBlock<Request>();
    var task = Task.Run(async () =>
    {
        while (await buffer.OutputAvailableAsync())
        {
            IList<Request> requests;
            buffer.TryReceiveAll(out requests);
            Calculate(requests.Last());
        }
    });
    

    用法:

    buffer.Post(new Request());
    buffer.Post(new Request());
    

    编辑:如果您没有Calculate 方法的任何输入或输出,您可以简单地使用boolean 作为开关。如果是真的,你可以关闭它并计算,如果它在Calculate 运行时再次变为真,那么再次计算:

    public bool _shouldCalculate;
    
    public void Producer()
    {
        _shouldCalculate = true;
    }
    
    public async Task Consumer()
    {
        while (true)
        {
            if (!_shouldCalculate)
            {
                await Task.Delay(1000);
            }
            else
            {
                _shouldCalculate = false;
                Calculate();
    
            }
        }
    }
    

    【讨论】:

    • Calculate 目前是一个没有参数的 void 方法,我似乎无法让你的代码工作。我猜 Linqpad 不是测试异步内容的最佳场所!
    • @DogEars TPL Dataflow 是一个 Nuget 库,您需要先添加它。但如果Calculate 没有输入或输出,则不需要队列。只需有一个布尔值,在请求时设置为 true,在计算时设置为 false。
    • 我已将 nuget 包添加到 LinqPad 中。我们需要确保如果正在进行计算。后续请求将被执行,因此我认为需要 1 的“队列” , 但它不需要存储任何东西.. 只是未完成请求的记录。
    • @DogEars 它可以简单地打开表示要发出请求的布尔值
    • @l3arnon 我们到了!目前计算是在调用线程上同步完成的。我需要避免 task.delay 我们可能有几十个具有计算方法的对象。我真的不想为每个人生成一个“消费者”吗?这不是让它变得更复杂。我已经更新了我的基本示例,按照您的建议进行了一些简化。干杯,耳朵。
    【解决方案2】:

    一次只需要 1 个的 BlockingCollection
    诀窍是如果集合中有任何项目则跳过

    我会接受 I3aron +1 的回答
    这(也许)是一个 BlockingCollection 解决方案

    public static void BC_AddTakeCompleteAdding()
    {
        using (BlockingCollection<int> bc = new BlockingCollection<int>(1))
        {
    
            // Spin up a Task to populate the BlockingCollection  
            using (Task t1 = Task.Factory.StartNew(() =>
            {
                for (int i = 0; i < 100; i++)
                {
                    if (bc.TryAdd(i))
                    {
                        Debug.WriteLine("  add  " + i.ToString());
                    }
                    else
                    {
                        Debug.WriteLine("  skip " + i.ToString());
                    }
    
                    Thread.Sleep(30);
                }
                bc.CompleteAdding();
            }))
            {
    
                // Spin up a Task to consume the BlockingCollection 
                using (Task t2 = Task.Factory.StartNew(() =>
                {
                    try
                    {
                        // Consume consume the BlockingCollection 
                        while (true)
                        {
                            Debug.WriteLine("take " + bc.Take());
                            Thread.Sleep(100);
                        }
                    }
                    catch (InvalidOperationException)
                    {
                        // An InvalidOperationException means that Take() was called on a completed collection
                        Console.WriteLine("That's All!");
                    }
                }))
    
                    Task.WaitAll(t1, t2);
            }
        }
    }
    

    【讨论】:

    • 这不是线程安全的。多个可以进入这个if (bc.Count == 0)并阻止bc.Add(i);
    • @I3arnon 好的,它认为我必须同意它不是线程安全的。即使在那种情况下,只有该线程会被阻塞,直到它被清除——其他线程仍然会跳过。你只会把那个弄坏。并且问题并未表明从多个线程中添加。
    • @I3arnon 好的,它确实说“请求来自多个线程”。我错过了那部分。
    【解决方案3】:

    这听起来像是典型的生产者-消费者。我建议查看BlockingCollection&lt;T&gt;。它是System.Collection.Concurrent 命名空间的一部分。最重要的是,您可以实现排队逻辑。

    您可以向BlockingCollection 提供任何内部结构来保存其数据,例如ConcurrentBag&lt;T&gt;ConcurrentQueue&lt;T&gt; 等。后者是使用的默认结构。

    【讨论】:

    • 您能详细解释一下吗?我无法理解 BlockingCollection 应该如何解决问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-01-12
    相关资源
    最近更新 更多