【发布时间】:2019-08-19 09:15:42
【问题描述】:
我正在使用消息队列中的消息并使用Task.Run() 并行处理它们。但是我想将消耗速度限制在某个最大线程数并且在线程数低于该线程数之前不要从消息队列中消耗。
假设我想要最多 100 个线程。在这种情况下,当达到 100 个线程时,它应该停止从消息队列中消费。当一个消息处理任务完成并且线程数下降到 99 时,它应该从队列中再消耗一条消息。
我尝试使用TransformBlock 来实现此目的,下面是用于演示的示例代码:
public partial class MainWindow : Window
{
object syncObj = new object();
int i = 0;
public MainWindow()
{
InitializeComponent();
}
private async Task<bool> ProcessMessage(string message)
{
await Task.Delay(5000);
lock (syncObj)
{
i++;
System.Diagnostics.Debug.WriteLine(i);
}
return true;
}
private async void Button_Click(object sender, RoutedEventArgs e)
{
var processor = new TransformBlock<string, bool>(
(str) => ProcessMessage(str),
new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 100 }
);
for(int i = 0; i < 1000; i++)
{
await processor.SendAsync("a");
}
}
}
限制并行任务的数量按预期工作,但所有消息都会立即发送到 TransformBlock,因此SendAsync 循环在任务处理之前结束。
只要线程数低于最大值,我希望它继续接受消息。允许并行,但在达到 100 时等待。
有没有办法使用 TransformBlock 来做到这一点,还是我应该求助于其他方法?
【问题讨论】:
-
为什么不使用 Parrallel.For 而不是 thr for 循环并在里面传递 MaxDegreeOfParallelism?还有任务不是线程,看看stackoverflow.com/questions/20806238/…
标签: c# multithreading