You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 22:00:17