【问题标题】:How to wait for a specific amount of time after creating a specific number of tasks\Threads?创建特定数量的任务\线程后如何等待特定的时间?
【发布时间】:2015-09-23 11:20:43
【问题描述】:

我有一个要求,我可以在一秒钟内访问一个 API 5 次。如果我必须发出总共 50 个请求,我想发出前 5 个请求并等待 1 秒,然后才能用另一批 5 个请求访问 API。我尝试使用线程池以及并行任务库 For\Foreach 循环和任务类,但我无法获得一个告诉我已创建 5 个任务的顺序计数器。 这是我正在尝试做的一个示例:

List<string> str = new List<string>();
for (int i = 0; i <= 100; i++)
{
    str.Add(i.ToString());
}

Parallel.ForEach(str, new ParallelOptions { MaxDegreeOfParallelism = 5 },
(value, pls, index) =>
{
    Console.WriteLine(value);// simulating method call
    if (index + 1 == 5)
    {
        // need the main thread to sleep so next batch is 
        Thread.Sleep(1000);
    }
});

【问题讨论】:

  • 您是否考虑过使用更好的方法?也许async/await
  • 我正在使用 .net 4.0。据我所知,async/await 可从 .net 4.5 获得。如果我错了,请纠正我。
  • 如果您有 VS 2012 或更高版本,您可以在 .net 4.0 中使用 async await。参考stackoverflow.com/questions/19421878/…
  • @SriramSakthivel 而不是建议可能不适用于 OP 的新工具,而是尝试使用代码解决问题。
  • @Gusdor 建议 OP 可能可能或可能不可用的工具是一个有效的建议,因为他可能根本不知道它们存在。例如,按照建议使用async-await,将使 OP 代码更简洁,并且可能更具可扩展性,而不是使用这种线程化方法。

标签: c# multithreading parallel-processing


【解决方案1】:

由于您使用的是 .NET 4.0(并且希望您至少使用 VS2012),因此您可以使用 Microsoft.Bcl.Async 来获得 async-await 功能。

一旦您这样做了,您就可以轻松地异步查询您的 API 端点(不需要额外的线程),并使用 AsyncSemaphore(参见下面的实现)来限制您同时执行的请求数。

例如:

public readonly AsyncSemaphore = new AsyncSemaphore(5);
public readonly HttpClient httpClient = new HttpClient();
public async Task<string> LimitedQueryAsync(string url)
{
    await semaphoreSlim.WaitAsync();
    try
    {
        var response = await httpClient.GetAsync(url);
        return response.Content.ReadAsStringAsync();
    }
    finally
    {
        semaphoreSlim.Release();
    }
}

现在可以这样查询了:

public async Task DoQueryStuffAsync()
{
    while (someCondition)
    {
        var results = await LimitedQueryAsync(url);

        // do stuff with results
        await Task.Delay(1000);
    }
}

编辑: 正如@ScottChamberlain 正确指出的那样,SemaphoreSlim 在.NET 4 中不可用。您可以改用AsyncSemaphore,如下所示:

public class AsyncSemaphore 
{ 
    private readonly static Task s_completed = Task.FromResult(true); 
    private readonly Queue<TaskCompletionSource<bool>> m_waiters = 
                                            new Queue<TaskCompletionSource<bool>>(); 
    private int m_currentCount; 

    public AsyncSemaphore(int initialCount)
    {
        if (initialCount < 0) 
        {
            throw new ArgumentOutOfRangeException("initialCount"); 
        }
        m_currentCount = initialCount; 
    }

    public Task WaitAsync() 
    { 
        lock (m_waiters) 
        { 
            if (m_currentCount > 0) 
            { 
                --m_currentCount; 
                return s_completed; 
            } 
            else 
            { 
                var waiter = new TaskCompletionSource<bool>(); 
                m_waiters.Enqueue(waiter); 
                return waiter.Task; 
            } 
        } 
    }

    public void Release() 
    { 
        TaskCompletionSource<bool> toRelease = null; 
        lock (m_waiters) 
        { 
            if (m_waiters.Count > 0) 
                toRelease = m_waiters.Dequeue(); 
            else 
                ++m_currentCount; 
        } 
        if (toRelease != null) 
            toRelease.SetResult(true); 
    }
}

【讨论】:

  • 鉴于Microsoft.Bcl.Async 可能不可用,您能否建议如何使用Task 完成此操作?
  • WaitAsync() 是否与Microsoft.Bcl.Async 一起添加(可能通过扩展方法)?它不是 4.0 中 SemaphoreSlim 的一部分
  • @ScottChamberlain 你是对的。我忘记了。修改了代码。
  • 我只是仔细检查了一下,它提供了System.Threading.Tasks.TaskExSystem.Net.DnsEx 两个类,并在AsyncExtensionsAsyncPlatformExtensionsAsyncPlatformExtensionsAwaitExtensions 中提供了各种扩展,但它们都没有处理信号量,其中大多数涉及使各种 IO 操作使用 TAP model 或扩展 TaskCancelationTokenSource 以添加 4.5 中添加的功能。
  • @ScottChamberlain 是的,我只记得我回答了another question,它问的正是:P
【解决方案2】:

如果已经限制为每秒 5 次,那么并行运行有多重要?这是尝试的不同视角(未经过编译测试)。这个想法是限制每个,而不是限制一个批次。

foreach(string value in values)
{
  const int alottedMilliseconds = 200;
  Stopwatch timer = Stopwatch.StartNew();

  // ...

  timer.Stop();
  int remainingMilliseconds = alottedMilliseconds - timer.ElapsedMilliseconds;
  if(remainingMilliseconds > 0)
  {
    // replace with something more precise/thread friendly as needed.
    Thread.Sleep(remainingMilliseconds);
  }
}

或者本着您最初要求的精神。使用扩展方法扩展您的解决方案,将您的列表分成 5 个块...

public static IEnumerable<List<T>> Partition<T>(this IList<T> source, Int32 size)
{
  for (int i = 0; i < Math.Ceiling(source.Count / (Double)size); i++)
  {
    yield return new List<T>(source.Skip(size * i).Take(size));
  }
}

利用此扩展在外循环中调用您的 Parallel.ForEach,然后在每个外循环结束时应用相同的计时器方法。像这样的...

foreach(IEnumerable<string> batch in str.Partitition(5))
{
  Stopwatch timer = Stopwatch.StartNew();

  Parallel.ForEach(
    batch, 
    new ParallelOptions { MaxDegreeOfParallelism = 5 },
    (value, pls, index) =>
    {
      Console.WriteLine(value);// simulating method call
    });

  timer.Stop();
  int remainingMilliseconds = 5000 - timer.ElapsedMilliseconds;
  if(remainingMilliseconds > 0)
  {
    // replace with something more precise/thread friendly as needed.
    Thread.Sleep(remainingMilliseconds);
  }
}

【讨论】:

  • 谢谢!!您的解决方案与我正在寻找的最接近。
【解决方案3】:

下面给出了两种方法。这两种方式,您都将获得所需的测试配置。 不仅代码简洁,而且不用加锁就可以实现。


1) 递归

您必须分批发出 50 个请求,每个请求 5 个。这意味着以 1 秒的间隔总共 10 批 5 个请求。定义实体,让:

  • HitAPI() 是一次调用 API 的线程安全方法;
  • InitiateBatch() 是启动一批 5 个线程来命中 API 的方法,

那么,示例实现可以是:

private void InitiateRecursiveHits(int batchCount)
{
    return InitiateBatch(batchCount);
}

只要用batchCount = 10 调用上面的方法,它就会调用下面的代码..

private void InitiateBatch(int batchNumber)
{
    if (batchNumber <= 0) return;
    var hitsPerBatch = 5;
    var thisBatchHits = new Task[hitsPerBatch];
    for (int taskNumber = 1; taskNumber <= hitsPerBatch; taskNumber++)
         thisBatchHits[taskNumber - 1] = Task.Run(HitAPI);
    Task.WaitAll(thisBatchHits);
    Thread.Sleep(1000); //To wait for 1 second before starting another batch of 5
    InitiateBatch(batchNumber - 1);
    return
}

2) 迭代

这比第一种方法更简单。只需以迭代的方式执行递归方法...

private void InitiateIterativeHits(int batchCount)
{
    if (batchCount <= 0) return;
    // It's good programming practice to leave your input variables intact so that 
    // they hold correct value throughout the execution
    int desiredRuns = batchCount;
    var hitsPerBatch = 5;
    while (desiredRuns-- > 0)
    {
        var thisBatchHits = new Task[hitsPerBatch];
        for (int taskNumber = 1; taskNumber <= hitsPerBatch; taskNumber++)
            thisBatchHits[taskNumber - 1] = Task.Run(HitAPI);
        Task.WaitAll(thisBatchHits);
        Thread.Sleep(1000); //To wait for 1 second before starting another batch of 5
    }
}

【讨论】:

    【解决方案4】:

    我会为此使用 Microsoft 的响应式框架 (NuGet "Rx-Main")。

    这就是它的样子:

    var query =
        Observable
            .Range(0, 100)
            .Buffer(5)
            .Zip(Observable.Interval(TimeSpan.FromSeconds(1.0)), (ns, i) => ns)
            .SelectMany(ns =>
                ns
                    .ToObservable()
                    .SelectMany(n =>
                        Observable
                            .Start(() =>
                            {
                                /* call here */
                                Console.WriteLine(n);
                                return n;
                            })));
    

    然后你会像这样处理结果:

    var subscription =
        query
            .Subscribe(x =>
            {
                /* handle result here */
            });
    

    如果您需要在请求自然完成之前停止请求,您只需致电subscription.Dispose();

    漂亮、干净、声明性。

    【讨论】:

      【解决方案5】:

      也许:

      while(true){
         for(int i = 0; i < 5; i++)
             Task.Run(() => { <API STUFF> });
         Thread.Sleep(1000);
      }
      

      我不确定一直这样调用 task.run 是否有效。

      【讨论】:

      • 注意这会阻塞主线程大约5秒
      • @MatiCicero 不,它会永远阻塞:)
      • @SriramSakthivel 你说得对!我仍然需要完全醒来:P
      • 显然这个答案并不意味着直接复制并粘贴到任何解决方案中,因为它确实会永远阻塞主线程。但是,可以在 while 循环中使用适当的条件在其他线程上运行它
      猜你喜欢
      • 1970-01-01
      • 2012-04-14
      • 1970-01-01
      • 2019-10-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多