【问题标题】:Task Parallel Library - Custom Task Schedulers任务并行库 - 自定义任务计划程序
【发布时间】:2012-03-20 21:54:13
【问题描述】:

我需要将 Web 服务请求发送到在线 api,我认为 Parallel Extensions 非常适合我的需求。

有问题的网络服务被设计为被重复调用,但有一种机制,如果你每秒调用超过一定数量,就会向你收费。我显然想尽量减少我的费用,所以想知道是否有人见过可以满足以下要求的 TaskScheduler:

  1. 限制每个时间跨度计划的任务数。我猜如果请求的数量超过了这个限制,那么它需要丢弃任务或可能阻塞? (停止积压的任务)
  2. 检测相同的请求是否已经在要执行的调度程序中但尚未执行,如果是,则不要将第二个任务排队,而是返回第一个任务。

人们是否觉得这些是任务调度程序应该处理的职责,还是我在找错树?如果您有其他选择,我愿意接受建议。

【问题讨论】:

  • 我认为至少 #2 不可能与 TaskScheduler 相关,因为它与 Tasks 打交道,并且无法从中获取这些信息。
  • 您对使用 C# 5/.Net 4.5 的解决方案感兴趣吗?
  • @svick 绝对,我很幸运被允许使用 c#5

标签: c# task-parallel-library parallel-extensions


【解决方案1】:

我同意其他人的观点,即 TPL 数据流听起来是一个很好的解决方案。

为了限制处理,您可以创建一个TransformBlock,它实际上不会以任何方式转换数据,如果它在前一个数据之后过早到达,它只会延迟它:

static IPropagatorBlock<T, T> CreateDelayBlock<T>(TimeSpan delay)
{
    DateTime lastItem = DateTime.MinValue;
    return new TransformBlock<T, T>(
        async x =>
                {
                    var waitTime = lastItem + delay - DateTime.UtcNow;
                    if (waitTime > TimeSpan.Zero)
                        await Task.Delay(waitTime);

                    lastItem = DateTime.UtcNow;

                    return x;
                },
        new ExecutionDataflowBlockOptions { BoundedCapacity = 1 });
}

然后创建一个产生数据的方法(例如从 0 开始的整数):

static async Task Producer(ITargetBlock<int> target)
{
    int i = 0;
    while (await target.SendAsync(i))
        i++;
}

它是异步写入的,因此如果目标块现在无法处理项目,它会等待。

然后写一个消费者方法:

static void Consumer(int i)
{
    Console.WriteLine(i);
}

最后,将它们链接在一起并启动它:

var delayBlock = CreateDelayBlock<int>(TimeSpan.FromMilliseconds(500));

var consumerBlock = new ActionBlock<int>(
    (Action<int>)Consumer,
    new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded });

delayBlock.LinkTo(consumerBlock, new DataflowLinkOptions { PropagateCompletion = true });

Task.WaitAll(Producer(delayBlock), consumerBlock.Completion);

这里,delayBlock 每 500 毫秒最多接受一个项目,Consumer() 方法可以并行运行多次。要完成处理,请致电delayBlock.Complete()

如果您想为每个 #2 添加一些缓存,您可以创建另一个 TransformBlock 在那里完成工作并将其链接到其他块。

【讨论】:

  • 宾果游戏,这正是我的想法。很高兴你提供了一个实际的实现,我只是没能找到时间。
  • 整洁。需要注意的一点是,如果您运行的是异步 CTP 而不是 .NET 4.5,则需要将 Task.Delay 更改为 TaskEx.Delay。
  • @DPeden,我认为您的意思是 .Net 4.5(目前处于测试阶段)。
【解决方案2】:

老实说,我会在更高的抽象级别上工作,并为此使用 TPL 数据流 API。唯一的问题是您需要编写一个自定义块,以您需要的速率限制请求,因为默认情况下,块是“贪婪的”,并且会尽可能快地处理。实现将是这样的:

  1. BufferBlock&lt;T&gt; 开头,这是您要发布到的逻辑块。
  2. BufferBlock&lt;T&gt; 链接到具有每秒请求数和限制逻辑知识的自定义块。
  3. 将自定义块从 2 链接到您的 ActionBlock&lt;T&gt;

我现在没有时间为 #2 编写自定义块,但是如果你还没有弄清楚,我会稍后再回来查看并尝试为你填写一个实现。

【讨论】:

  • 由于我根本没有研究过 TPL 数据流,非常感谢您的帮助
【解决方案3】:

我没有太多使用 RX,但是 AFAICT Observable.Window 方法可以很好地解决这个问题。

http://msdn.microsoft.com/en-us/library/system.reactive.linq.observable.window(VS.103).aspx

它似乎比 Throttle 更合适,它似乎会丢弃元素,我猜这不是你想要的

【讨论】:

    【解决方案4】:

    如果您需要按时间节流,请查看Quartz.net。它可以促进一致的轮询。如果您关心所有请求,则应考虑使用某种排队机制。 MSMQ 可能是正确的解决方案,但如果您想扩大规模并使用像 NServiceBusRabbitMQ 这样的 ESB,则有许多特定的实现。

    更新:

    在这种情况下,如果您可以利用 CTP,TPL 数据流是您的首选解决方案。一个受限制的 BufferBlock 就是解决方案。

    这个例子来自documentation provided by Microsoft

    // Hand-off through a bounded BufferBlock<T>
    private static BufferBlock<int> m_buffer = new BufferBlock<int>(
        new DataflowBlockOptions { BoundedCapacity = 10 });
    
    // Producer
    private static async void Producer()
    {
        while(true)
        {
            await m_buffer.SendAsync(Produce());
        }
    }
    
    // Consumer
    private static async Task Consumer()
    {
        while(true)
        {
            Process(await m_buffer.ReceiveAsync());
        }
    }
    
    // Start the Producer and Consumer
    private static async Task Run()
    {
        await Task.WhenAll(Producer(), Consumer());
    }
    

    更新:

    查看 RX 的Observable.Throttle

    【讨论】:

    • 我想保持一切正常,因此排除了排队作为解决方案。我也不认为 Quartz 是正确的解决方案,调度不是问题,它限制了请求的数量并且对我的调用有点聪明。感谢您的建议
    • @DarenFox 这究竟是如何阻止每秒超过 N 个请求的?看起来您只限制了 10 个并发未完成的调用,但显然如果调用速度超过 10 毫秒,则可能超过每秒限制。 FWIW,我的回答是 TPL 数据流可能也是最好的方法,但从技术上讲,您需要一个临时的 BufferBlock 实现。
    • @DarenFox 这也是我关心的问题,但根据 DarrenFox 的最后评论,它似乎被丢弃了。尽管如此,我认为RX可能是解决方案。我用另一个选项再次更新了我的答案。
    • 据我了解 RX Throttle 无济于事,因为它是 observable 上的选择器。这意味着仍然需要向 Web 服务发出请求,但在该时间段内没有更多请求之前不会在 observable 中返回。实际上,我使用 Observable.Timer 实现了我的要求的部分实现,但它变得越来越笨重,这就是为什么我退后一步来评估其他选项
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-10-23
    • 2011-07-23
    • 1970-01-01
    • 2020-07-25
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多