【发布时间】:2018-10-11 17:12:38
【问题描述】:
我们在使用 TPL 数据流库时需要一个请求/响应模式。我们的问题是我们有一个调用依赖服务的 .NET 核心 API。依赖服务限制并发请求。我们的 API 不限制并发请求;因此,我们一次可以收到数千个请求。在这种情况下,依赖服务将在达到其限制后拒绝请求。因此,我们实现了BufferBlock<T> 和TransformBlock<TIn, TOut>。性能稳定,效果很好。我们测试了我们的 API 前端,有 1000 个用户发出 100 个请求/秒,0 个问题。缓冲区块缓冲请求,转换块并行执行我们所需数量的请求。依赖服务接收我们的请求并做出响应。我们在转换块操作中返回该响应,一切都很好。我们的问题是缓冲区块和转换块断开连接,这意味着请求/响应不同步。我们遇到了一个请求将收到另一个请求者的响应的问题(请参阅下面的代码)。
具体到下面的代码,我们的问题出在GetContent方法上。该方法是从我们 API 中的服务层调用的,该服务层最终是从我们的控制器调用的。下面的代码和服务层是单例的。缓冲区的SendAsync 与转换块ReceiveAsync 断开连接,因此返回任意响应而不一定是发出的请求。
所以,我们的问题是:有没有办法使用数据流块来关联请求/响应?最终目标是请求进入我们的 API,发送到依赖服务,然后返回给客户端。我们的数据流实现代码如下。
public class HttpClientWrapper : IHttpClientManager
{
private readonly IConfiguration _configuration;
private readonly ITokenService _tokenService;
private HttpClient _client;
private BufferBlock<string> _bufferBlock;
private TransformBlock<string, JObject> _actionBlock;
public HttpClientWrapper(IConfiguration configuration, ITokenService tokenService)
{
_configuration = configuration;
_tokenService = tokenService;
_bufferBlock = new BufferBlock<string>();
var executionDataFlowBlockOptions = new ExecutionDataflowBlockOptions
{
MaxDegreeOfParallelism = 10
};
var dataFlowLinkOptions = new DataflowLinkOptions
{
PropagateCompletion = true
};
_actionBlock = new TransformBlock<string, JObject>(t => ProcessRequest(t),
executionDataFlowBlockOptions);
_bufferBlock.LinkTo(_actionBlock, dataFlowLinkOptions);
}
public void Connect()
{
_client = new HttpClient();
_client.DefaultRequestHeaders.Add("x-ms-client-application-name",
"ourappname");
}
public async Task<JObject> GetContent(string request)
{
await _bufferBlock.SendAsync(request);
var result = await _actionBlock.ReceiveAsync();
return result;
}
private async Task<JObject> ProcessRequest(string request)
{
if (_client == null)
{
Connect();
}
try
{
var accessToken = await _tokenService.GetTokenAsync(_configuration);
var httpRequestMessage = new HttpRequestMessage(HttpMethod.Post,
new Uri($"https://{_configuration.Uri}"));
// add the headers
httpRequestMessage.Headers.Add("Authorization", $"Bearer {accessToken}");
// add the request body
httpRequestMessage.Content = new StringContent(request, Encoding.UTF8,
"application/json");
var postRequest = await _client.SendAsync(httpRequestMessage);
var response = await postRequest.Content.ReadAsStringAsync();
return JsonConvert.DeserializeObject<JObject>(response);
}
catch (Exception ex)
{
// log error
return new JObject();
}
}
}
【问题讨论】:
标签: c# .net task-parallel-library tpl-dataflow