【问题标题】:Set a limit to parallel tasks received by azure DownloadToStreamAsync对 azure DownloadToStreamAsync 接收的并行任务设置限制
【发布时间】:2014-02-03 11:01:01
【问题描述】:

我有一堆文件(大约 10k)需要从 Windows azure 存储下载。为了让他们并行下载而不是一次下载一个,我使用了 blob DownloadToStreamAsync 方法,该方法返回一个 Task 对象。然后,我使用将流保存到文件的方法设置任务 ContinueWith。

代码如下:

foreach (var File in ServerFiles)
{
    string sFileName = File.Uri.LocalPath.ToString();
    CloudBlockBlob oBlob = BiActionscontainer.GetBlockBlobReference(sFileName.Replace("/" + Container + "/", ""));

    MemoryStream ms = new MemoryStream();
    BlobRequestOptions f = new BlobRequestOptions();
    Task downloadTask = oBlob.DownloadToStreamAsync(ms);

    downloadTask.ContinueWith((Task task) =>
    {
         ms.Position = 0;
         lock(lockObject)
         {
              using (FileStream file = new FileStream(ResultPath, FileMode.Append, FileAccess.Write))
              {
                   byte[] bytes = ms.ToArray();
                   file.Write(bytes, 0, bytes.Length);
              }
         }
         ms.Dispose();
    });
}

此代码是在我们的一台服务器(不是在 azure 上)运行的工具的一部分 - windows 2003 服务器。问题是在该服务器上我得到“操作已超时。Windows 2003 标准上的 Microsoft.WindowsAzure.Storage”,所以我认为可能是很多文件同时发出请求并阻塞了带宽.

所以我想知道,在我从第三方库获取 Task 对象的这种情况下,如何限制一次运行的并行数?并且仍然将其余的任务排队?

【问题讨论】:

  • 如果我没记错的话,你的系统因为下面的代码行foreach (var File in ServerFiles)而被破坏了。您需要以某种方式在此处节流。

标签: c# azure task-parallel-library


【解决方案1】:

您可以为此使用SemaphoreSlim。将其设置为您想要的并发请求数,然后在开始每个请求之前使用await WaitAsync(),在每个请求完成后使用Release(),最后等待剩余的任务。

封装在一个辅助方法中,它可能看起来像这样:

public static async Task ForEachAsync<T>(
    this IEnumerable<T> items, Func<T, Task> action, int maxDegreeOfParallelism)
{
    var semaphore = new SemaphoreSlim(maxDegreeOfParallelism);

    var tasks = new List<Task>();

    foreach (var item in items)
    {
        await semaphore.WaitAsync();

        Func<T, Task> loopAction = async x =>
        {
            await action(x);
            semaphore.Release();
        };

        tasks.Add(loopAction(item));
    }

    await Task.WhenAll(tasks);
}

用法(对您的代码进行了一些更改,主要是为了简化它并使其更加异步):

ServerFiles.ForEachAsync(async file =>
{
    string sFileName = File.Uri.LocalPath.ToString();
    CloudBlockBlob oBlob = BiActionscontainer.GetBlockBlobReference(sFileName.Replace("/" + Container + "/", ""));

    var ms = new MemoryStream();
    BlobRequestOptions f = new BlobRequestOptions();
    await oBlob.DownloadToStreamAsync(ms);

    ms.Position = 0;
    lock (lockObject)
    {
         using (var file = new FileStream(ResultPath, FileMode.Append, FileAccess.Write))
         {
              await ms.CopyToAsync(file);
         }
    }
});

另一种实现是使用来自 TPL Dataflow 的 ActionBlock。它知道这里需要的一切,你只需要设置它:

public static Task ForEachAsync<T>(
    this IEnumerable<T> items, Func<T, Task> action, int maxDegreeOfParallelism)
{
    var block = new ActionBlock<T>(
        action,
        new ExecutionDataflowBlockOptions
        {
            MaxDegreeOfParallelism = maxDegreeOfParallelism
        });

    foreach (var item in items)
    {
        block.Post(item);
    }

    block.Complete();
    return block.Completion;
}

【讨论】:

    猜你喜欢
    • 2014-02-20
    • 1970-01-01
    • 2012-07-25
    • 2015-04-16
    • 1970-01-01
    • 2020-07-28
    • 2019-08-11
    • 2022-10-15
    • 1970-01-01
    相关资源
    最近更新 更多