Ocelot自定义负载均衡多线程批量请求7/3分配异常排查
问题:多线程下Ocelot网关批量请求7:3负载分配失效
我实现了一个向Ocelot网关发送批量请求的Worker,网关使用自定义负载均衡器MyRoundRobin,期望将每个批量请求的70%分配到HostAndPort1,30%分配到另一节点,最终实现总请求7:3的分发比例。但当前仅单线程下正常,多线程场景无法按预期分配。
1. Worker批量请求发送代码
public async ValueTask<long> DoReadAsync() { OnProcess = true; var count = 0; try { var sends = (await _unitOfWork.SendRepository.Where(x => x.sndPost == false, x => new Send { sndID = x.sndID, SndBody = x.SndBody, SndTo = x.SndTo, sndFrom = x.sndFrom, SndFarsi = x.SndFarsi, sndMsgClass = x.sndMsgClass, sndUDH = x.sndUDH }, c => c.OrderBy(o => o.SndPriority), BaseConfig.BatchRead)).ToList(); count = sends.Count; if (count == 0) { await Task.Delay(_settings.RefreshMillisecond); OnProcess = false; return count; } var split = BaseConfig.Split(sends); var restTasks = split .Select(items => _restService.DoPostAsync(items.ToList())) .ToList(); var updateTasks = new List<Task>(); while (restTasks.Any()) { var task = await Task.WhenAny(restTasks); //task.ThrowExceptionIfTaskIsFaulted(); var item = await task.ConfigureAwait(false); if (item is { IsSuccessStatusCode: true }) { var content = await item.Content.ReadAsStringAsync(); var items = JsonConvert.DeserializeObject<List<ResponseDto>>(content); var itemsSends = items.Select(_mapper.Map<Send>).ToList(); if (itemsSends.Any()) { var updateTask = _unitOfWork.SendRepository.BulkUpdateForwardOcelotAsync(itemsSends); updateTasks.Add(updateTask); } } restTasks.Remove(task); } await Task.WhenAll(updateTasks).ConfigureAwait(false); Completed = true; OnProcess = false; } catch (Exception ex) { _logger.LogError(ex.Message); OnProcess = false; } return count; }
2. Ocelot网关BatchMiddleware中间件代码
public class BatchMiddleware : OcelotMiddleware { private readonly RequestDelegate _next; private bool isRahyabRG = true; private int remainedBatch = 0; public BatchMiddleware( RequestDelegate next, IConfiguration configuration, IOcelotLoggerFactory loggerFactory) : base(loggerFactory.CreateLogger<BatchMiddleware>()) { _next = next; } public async Task Invoke(HttpContext httpContext) { var request = httpContext.Request; var batchRequests = await request.DeserializeArrayAsync<RequestDto>(); var batchRequestCount = batchRequests.Count; var RGCount = (int)Math.Floor(70 * batchRequestCount / 100.0); if (isRahyabRG) { var rgRequests = batchRequests.Take(RGCount).ToList(); var requestBody = JsonConvert.SerializeObject(rgRequests); request.Body = new MemoryStream(Encoding.UTF8.GetBytes(requestBody)); isRahyabRG = false; remainedBatch = batchRequestCount - RGCount; httpContext.Session.SetString("remainedBatchKey", remainedBatch.ToString()); } else { var remainedBatchKey = httpContext.Session.GetString("remainedBatchKey"); var pmRequests = new List<RequestDto>(); if (remainedBatchKey != null) { pmRequests = batchRequests.Take(int.Parse(remainedBatchKey)).ToList(); } var requestBody = JsonConvert.SerializeObject(pmRequests); request.Body = new MemoryStream(Encoding.UTF8.GetBytes(requestBody)); isRahyabRG = true; } await _next.Invoke(httpContext); } }
3. MyRoundRobin自定义负载均衡器代码
public class MyRoundRobin : ILoadBalancer { private readonly Func<Task<List<Service>>> _services; private readonly object _lock = new(); private int _last; public MyRoundRobin(Func<Task<List<Service>>> services, IConfiguration configuration) { _services = services; } public async Task<Response<ServiceHostAndPort>> Lease(HttpContext httpContext) { var services = await _services(); lock (_lock) { if (_last >= services.Count) _last = 0; var next = services[_last++]; return new OkResponse<ServiceHostAndPort>(next.HostAndPort); } } public void Release(ServiceHostAndPort hostAndPort) { } }
4. ocelot.json配置
{ "Routes": [ { "DownstreamPathTemplate": "/api/Forward", "DownstreamScheme": "http", "DownstreamHostAndPorts": [ { "Host": "localhost", "Port": 51003 }, { "Host": "localhost", "Port": 32667 } ], "UpstreamPathTemplate": "/", "UpstreamHttpMethod": [ "POST" ], "LoadBalancerOptions": { "Type": "MyRoundRobin" } } ] }
内容的提问来源于stack exchange,提问作者Mansoure
相关产品推荐
相关产品推荐

