【问题标题】:Send parallel requests but only one per host with HttpClient and Polly to gracefully handle 429 responses使用 HttpClient 和 Polly 发送并行请求,但每个主机只能发送一个请求,以优雅地处理 429 响应
【发布时间】:2019-07-13 20:53:29
【问题描述】:

简介:

我正在构建一个单节点网络爬虫来简单地验证 .NET Core 控制台应用程序中的 URL 是 200 OK。我在不同的主机上有一组 URL,我用HttpClient 向它们发送请求。我对使用 Polly 和 TPL 数据流还很陌生。

要求:

  1. 我希望支持同时发送多个 HTTP 请求 可配置MaxDegreeOfParallelism。
  2. 我想将任何给定主机的并行请求数限制为 1(或可配置)。这是为了优雅地使用 Polly 策略处理每个主机的 429 TooManyRequests 响应。或者,我可以使用断路器在收到一个429 响应时取消对同一主机的并发请求,然后一次一个地处理该特定主机?
  3. 我完全可以完全不使用 TPL 数据流,而是可能使用 Polly Bulkhead 或其他一些机制来限制并行请求,但我不确定为了实现要求 2,该配置会是什么样子。

当前实施:

我当前的实现工作,除了我经常看到我将有x 并行请求到同一主机返回429 大约在同一时间...然后,他们都暂停重试策略.. . 然后,他们都在同一时间再次猛击同一个主机,经常仍然收到429s。即使我将同一主机的多个实例均匀地分布在整个队列中,我的 URL 集合也会被一些特定主机过度加权,这些主机最终仍会开始生成 429s。

收到429后,我想我只想向该主机发送一个并发请求,以尊重远程主机并追求200s。

验证器方法:

public async Task<int> GetValidCount(IEnumerable<Uri> urls, CancellationToken cancellationToken)
{
    var validator = new TransformBlock<Uri, bool>(
        async u => (await _httpClient.GetAsync(u, HttpCompletionOption.ResponseHeadersRead, cancellationToken)).IsSuccessStatusCode,
        new ExecutionDataflowBlockOptions {MaxDegreeOfParallelism = MaxDegreeOfParallelism}
    );
    foreach (var url in urls)
        await validator.SendAsync(url, cancellationToken);
    validator.Complete();
    var validUrlCount = 0;
    while (await validator.OutputAvailableAsync(cancellationToken))
    {
        if(await validator.ReceiveAsync(cancellationToken))
            validUrlCount++;
    }
    await validator.Completion;
    return validUrlCount;
}

Polly 策略应用于上述GetValidCount() 中使用的 HttpClient 实例。

IAsyncPolicy<HttpResponseMessage> waitAndRetryTooManyRequests = Policy
    .HandleResult<HttpResponseMessage>(r => r.StatusCode == HttpStatusCode.TooManyRequests)
    .WaitAndRetryAsync(3,
        (retryCount, response, context) =>
            response.Result?.Headers.RetryAfter.Delta ?? TimeSpan.FromMilliseconds(120),
        async (response, timespan, retryCount, context) =>
        {
            // log stuff
        });

问题:

如何修改或替换此解决方案以满足要求 #2?

【问题讨论】:

    标签: c# .net-core web-crawler tpl-dataflow polly


    【解决方案1】:

    我会尝试引入某种标志LimitedMode 来检测此特定客户端是否以受限模式进入。下面我声明了两个策略 - 一个简单的重试策略只是为了捕获 TooManyRequests 并设置标志。第二个策略是开箱即用的BulkHead 策略。

        public void ConfigureServices(IServiceCollection services)
        {
            /* other configuration */
    
            var registry = services.AddPolicyRegistry();
    
            var catchPolicy = Policy.HandleResult<HttpResponseMessage>(r =>
                {
                    LimitedMode = r.StatusCode == HttpStatusCode.TooManyRequests;
                    return false;
                })
                .WaitAndRetryAsync(1, i => TimeSpan.FromSeconds(3)); 
    
            var bulkHead = Policy.BulkheadAsync<HttpResponseMessage>(1, 10, OnBulkheadRejectedAsync);
    
            registry.Add("catchPolicy", catchPolicy);
            registry.Add("bulkHead", bulkHead);
    
            services.AddHttpClient<CrapyWeatherApiClient>((client) =>
            {
                client.BaseAddress = new Uri("hosturl");
            }).AddPolicyHandlerFromRegistry(PolicySelector);
        }
    

    然后您可能希望使用PolicySelector 机制动态决定应用哪个策略:如果受限模式处于活动状态 - 使用 catch 429 策略包装块头策略。如果收到成功状态代码 - 切换回没有隔板的常规模式。

        private IAsyncPolicy<HttpResponseMessage> PolicySelector(IReadOnlyPolicyRegistry<string> registry, HttpRequestMessage request)
        {
            var catchPolicy = registry.Get<IAsyncPolicy<HttpResponseMessage>>("catchPolicy");
            var bulkHead = registry.Get<IAsyncPolicy<HttpResponseMessage>>("bulkHead");
            if (LimitedMode)
            {
                return catchPolicy.WrapAsync(bulkHead);
            }
    
            return catchPolicy;
        }        
    

    【讨论】:

      【解决方案2】:

      这是一个创建TransformBlock 的方法,它可以防止具有相同密钥的消息并发执行。每条消息的密钥是通过调用提供的keySelector 函数获得的。具有相同密钥的消息相互顺序处理(而不是并行处理)。密钥也作为参数传递给transform 函数,因为它在某些情况下很有用。

      public static TransformBlock<TInput, TOutput>
          CreateExclusivePerKeyTransformBlock<TInput, TKey, TOutput>(
          Func<TInput, TKey, Task<TOutput>> transform,
          ExecutionDataflowBlockOptions dataflowBlockOptions,
          Func<TInput, TKey> keySelector,
          IEqualityComparer<TKey> keyComparer = null)
      {
          if (transform == null) throw new ArgumentNullException(nameof(transform));
          if (keySelector == null) throw new ArgumentNullException(nameof(keySelector));
          if (dataflowBlockOptions == null)
              throw new ArgumentNullException(nameof(dataflowBlockOptions));
          keyComparer = keyComparer ?? EqualityComparer<TKey>.Default;
      
          var internalCTS = CancellationTokenSource
              .CreateLinkedTokenSource(dataflowBlockOptions.CancellationToken);
      
          var maxDOP = dataflowBlockOptions.MaxDegreeOfParallelism;
          var taskScheduler = dataflowBlockOptions.TaskScheduler;
      
          var perKeySemaphores = new ConcurrentDictionary<TKey, SemaphoreSlim>(
              keyComparer);
      
          SemaphoreSlim maxDopSemaphore;
          if (maxDOP == DataflowBlockOptions.Unbounded)
          {
              maxDopSemaphore = new SemaphoreSlim(Int32.MaxValue);
          }
          else
          {
              maxDopSemaphore = new SemaphoreSlim(maxDOP, maxDOP);
      
              // The degree of parallelism is controlled by the semaphore
              dataflowBlockOptions.MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded;
      
              // Use a limited-concurrency scheduler for preserving the processing order
              dataflowBlockOptions.TaskScheduler = new ConcurrentExclusiveSchedulerPair(
                  taskScheduler, maxDOP).ConcurrentScheduler;
          }
      
          var block = new TransformBlock<TInput, TOutput>(async item =>
          {
              var key = keySelector(item);
              var perKeySemaphore = perKeySemaphores
                  .GetOrAdd(key, _ => new SemaphoreSlim(1, 1));
      
              // Continue on captured context before invoking the transform
              await perKeySemaphore.WaitAsync(internalCTS.Token);
              try
              {
                  await maxDopSemaphore.WaitAsync(internalCTS.Token);
                  try
                  {
                      return await transform(item, key).ConfigureAwait(false);
                  }
                  catch (Exception ex) when (!(ex is OperationCanceledException))
                  {
                      internalCTS.Cancel(); // The block has failed
                      throw;
                  }
                  finally
                  {
                      maxDopSemaphore.Release();
                  }
              }
              finally
              {
                  perKeySemaphore.Release();
              }
          }, dataflowBlockOptions);
      
          dataflowBlockOptions.MaxDegreeOfParallelism = maxDOP; // Restore initial value
          dataflowBlockOptions.TaskScheduler = taskScheduler; // Restore initial value
          return block;
      }
      

      使用示例:

      var validator = CreateExclusivePerKeyTransformBlock<Uri, string, bool>(
          async (uri, host) =>
          {
              return (await _httpClient.GetAsync(uri, HttpCompletionOption
                  .ResponseHeadersRead, token)).IsSuccessStatusCode;
          },
          new ExecutionDataflowBlockOptions
          {
              MaxDegreeOfParallelism = 30,
              CancellationToken = token,
          },
          keySelector: uri => uri.Host,
          keyComparer: StringComparer.OrdinalIgnoreCase);
      

      支持所有execution options(MaxDegreeOfParallelism、BoundedCapacity、CancellationToken、EnsureOrdered 等)。

      下面是接受同步委托的CreateExclusivePerKeyTransformBlock 的重载,以及另一个返回ActionBlock 而不是TransformBlock 的方法+重载,具有相同的行为。

      public static TransformBlock<TInput, TOutput>
          CreateExclusivePerKeyTransformBlock<TInput, TKey, TOutput>(
          Func<TInput, TKey, TOutput> transform,
          ExecutionDataflowBlockOptions dataflowBlockOptions,
          Func<TInput, TKey> keySelector,
          IEqualityComparer<TKey> keyComparer = null)
      {
          if (transform == null) throw new ArgumentNullException(nameof(transform));
          return CreateExclusivePerKeyTransformBlock(
              (item, key) => Task.FromResult(transform(item, key)),
              dataflowBlockOptions, keySelector, keyComparer);
      }
      
      // An ITargetBlock is similar to an ActionBlock
      public static ITargetBlock<TInput>
          CreateExclusivePerKeyActionBlock<TInput, TKey>(
          Func<TInput, TKey, Task> action,
          ExecutionDataflowBlockOptions dataflowBlockOptions,
          Func<TInput, TKey> keySelector,
          IEqualityComparer<TKey> keyComparer = null)
      {
          if (action == null) throw new ArgumentNullException(nameof(action));
          var block = CreateExclusivePerKeyTransformBlock(async (item, key) =>
              { await action(item, key).ConfigureAwait(false); return (object)null; },
              dataflowBlockOptions, keySelector, keyComparer);
          block.LinkTo(DataflowBlock.NullTarget<object>());
          return block;
      }
      
      public static ITargetBlock<TInput>
          CreateExclusivePerKeyActionBlock<TInput, TKey>(
          Action<TInput, TKey> action,
          ExecutionDataflowBlockOptions dataflowBlockOptions,
          Func<TInput, TKey> keySelector,
          IEqualityComparer<TKey> keyComparer = null)
      {
          if (action == null) throw new ArgumentNullException(nameof(action));
          return CreateExclusivePerKeyActionBlock(
              (item, key) => { action(item, key); return Task.CompletedTask; },
              dataflowBlockOptions, keySelector, keyComparer);
      }
      

      警告:这个类为每个键分配一个SemaphoreSlim,并保持对它的引用,直到类实例最终被垃圾回收。如果不同键的数量很大,这可能是一个问题。有一个分配较少的异步锁here 的实现,它在内部仅存储当前正在使用的SemaphoreSlims(加上一小部分可以重用的已释放SemaphoreSlims),它可以替换@此实现使用 987654341@。

      【讨论】:

      • 注意: 上述实现没有高效的输入队列。如果你用数百万个项目喂这个块,内存使用会爆炸。
      猜你喜欢
      • 1970-01-01
      • 2015-07-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多