【问题标题】:How to stop thread starvation when using ThreadPool.QueueUserWorkItem?使用 ThreadPool.QueueUserWorkItem 时如何停止线程饥饿?
【发布时间】:2021-06-11 05:33:15
【问题描述】:

我有一个 dot net core 5 控制台应用程序,它每分钟处理来自 rabbitmq 的大约 100,000 多条消息

当从 rabbitmq 接收到消息时,一个线程会关闭并处理一些数字,但是其中一个操作是调用外部 API 以获取有关其位置的信息。 当这个外部 API 服务速度变慢并且响应时间增加时,我看到 Windows 任务管理器上的线程不足和线程计数可能会达到 1000 个,并且应用程序基本上什么都不做

当应用程序加载时,主线程建立与rabbitmq的连接并订阅到达rabbitmq的新消息,每次消息到达时,我的控制台应用程序都会消耗每条消息,并启动一个线程池项目并继续获取新的rabbitmq消息

private void Consumer_Received(object sender, BasicDeliverEventArgs deliveryArgs)
{
    var data = Encoding.UTF8.GetString(deliveryArgs.Body.ToArray());
    ThreadPool.QueueUserWorkItem(new WaitCallback(StartProcessing), data);
}

如果我放置一个断点,这个 void 会一直被命中,一个新的线程池进程调用 StartProcessing void,这是 CPU 运行和外部 api 调用发生的地方

public void StartProcessing(object xdata)
{
    //1. crunch cpu
    
    //2. call external API
    
}

处理 cpu 的每条消息大约需要 100 毫秒,但正常情况下外部 API 需要 80-500 毫秒,但是当出现问题(可能是网络)时,最多可能需要 10 秒才能响应 1 个请求,这是应用程序开始崩溃的时候。

我的问题是围绕这个实现以及如何停止线程饥饿。 这是一个高吞吐量的多线程应用程序,它需要处理尽可能多的消息。 当外部 API 响应缓慢及其不断的上下文切换线程时,应用程序需要缓解背压。

使用 ThreadPool.QueueUserWorkItem 是正确的实现还是应该使用 Async await 等?

我也愿意听取这是否是一个糟糕的实现,以及我是否应该为此使用另一种模式。

///////////////////////////// 更新 1 ////////////////////////////////////p>

所以我更改了代码以使用异步任务,并且从 rabbitmq 获取消息的速度非常慢 旧代码在几秒钟内收到所有消息(200,000 条),新代码在几分钟内收到大约 1,000 条

新代码是

private void Consumer_Received(object sender, BasicDeliverEventArgs deliveryArgs)
{
    StartProcessing(deliveryArgs.Body.ToArray()).ConfigureAwait(false);
}
    
public static async Task<bool> StartProcessing(ReadOnlyMemory<byte> data)
{
    await Task.Run(() =>
    {
        ReadOnlySpan<byte> xdata = data.Span; //defensiveCopy of in memory pointer

        //do stuff
        
    }).ConfigureAwait(false);

return true;
}

我的实现有问题吗? “StartProcessing”代码应该被触发并忘记了,因为主线程应该继续到rabbitmq中的下一条消息 我似乎在等待消息处理后再继续 ////////////////////////////////////p>

【问题讨论】:

  • 这能回答你的问题吗? c# Threadpool - limit number of threads
  • 感谢@Sinatr,这很有趣,但想知道是否有任何更新的技术,因为那是 9 年前的帖子
  • 1) 外部 API 服务是否提供异步端点? 2) 是否有关于可以调用外部 API 服务的最大频率或并发性的任何记录限制? 3) 如果外部 API 服务无法跟上步伐,并且响应时间越来越长,您的应用程序的理想行为是什么?
  • 我们正在将外部 api 移动到本地服务器,但是它是一个数据库查找,所以它仍然很慢,它现在可以处理多个请求,并且在本地时,如果api 跟不上,那就是当线程变高时,我希望异步等待的东西可以为我处理这个问题。
  • 你应该尝试优化数据库查找。如果瓶颈在服务器,优化客户端不太可能带来任何好处。

标签: c# multithreading .net-core rabbitmq


【解决方案1】:

听起来这就是异步函数的确切场景。

如果您使用 CPU,则使用后台线程会对您有所帮助,但仅限于您拥有的硬件线程数。

但听起来您主要是在阻止网络 IO。使用一个在某种 IO 响应之前被阻塞的线程是非常浪费的,因为每个线程都会消耗一些资源。而且很容易导致线程池最大化等问题。

到现在为止,.Net 和许多库已经更新,为 IO 提供真正的异步函数。这释放线程来做其他事情而不是阻塞,当 IO 完成时,它将把剩余的工作安排在一个新的后台线程上。并且使用 async/await 可以让您或多或少地像编写常规同步代码一样编写代码,让编译器将其重写为状态机来处理维护状态的复杂问题。理想情况下,您需要的线程数不应超过您拥有的硬件线程数,因为每个线程都应该执行实际工作。

请记住,仅仅因为有一个异步方法返回一个任务,并不一定意味着它是真正的异步的。一些基类/接口,如流,已使用异步版本进行了扩展。而一些库供应商并没有提供实际的异步实现,只是封装了同步方法,没有提供真正的好处。

例如:

private async void Consumer_Received(...)
{
    try{
        var result = await Task.Run(()=> MyCpuBoundWork());    
        await MyNetworkCall(result);            
    }   
    catch{
         // handle exceptions
    } 
}

收到消息后,这将使用另一个后台线程来完成 CPU 密集型工作。我不确定 rabbitMq 是如何生成消息的,只有当它对所有消息使用单个线程时才需要 Task.Run 部分。在 CPU 绑定完成后,它将继续进行网络调用。

【讨论】:

  • 感谢@JonasH,我已经阅读了很多关于 IO 绑定和 CPU 绑定工作的内容,您应该为每个工作使用不同的 task.run/async await 调用,问题是当我从 rabbitmq 收到消息时,我需要在那个线程上同时做 CPU 和 IO 工作,所以不确定这是不好的做法?
  • @quade 添加示例
  • 谢谢,我已经用最新的信息/问题更新了我的原始帖子
猜你喜欢
  • 1970-01-01
  • 2012-07-26
  • 1970-01-01
  • 1970-01-01
  • 2013-04-01
  • 1970-01-01
  • 1970-01-01
  • 2016-06-23
  • 2017-12-15
相关资源
最近更新 更多