【问题标题】:Throttling Message Queue Consumption While Using Parallellism使用并行性时限制消息队列消耗
【发布时间】: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 来做到这一点,还是我应该求助于其他方法?

【问题讨论】:

标签: c# multithreading


【解决方案1】:

数据流块具有输入缓冲区。此输入缓冲区充当队列。

如果您想将消息保留在您的自己的队列中,您可以通过限制数据流块愿意接收的项目数来做一些接近您想要的事情:

var processor = new TransformBlock<string, bool>(
    (str) => ProcessMessage(str),
    new ExecutionDataflowBlockOptions
    {
      BoundedCapacity = 100,
      MaxDegreeOfParallelism = 100,
    }
);

请注意,BoundedCapacity 包括块正在处理的项目。由于BoundedCapacity == MaxDegreeOfParallelism,这实际上关闭了数据流块的队列。

所以 SendAsync 循环在任务处理之前结束。

当有(最多)100 个任务要处理时,它仍然会结束。如果您想等到所有项目都完成处理,请致电Complete() 和await Completed。

【讨论】:

  • 感谢您指出 BoundedCapacity。我会尽快尝试。由于我将在一段时间循环中不断地收听来自队列的消息,我认为我不需要调用 Complete()。我将依赖 BoundedCapacity 并等待 SendAsync。
  • 当我添加 BoundedCapacity=100 时,Debug.Writeline 在 100 次迭代后不起作用。
  • @hakanviv 您需要从TransformBlock 中获取结果。将其链接到另一个块或添加消费者。
  • 我有点困惑。我对 ProcessMessage 方法的虚拟结果并不感兴趣,我只是希望它输出到调试输出窗口。你的意思是它不会继续,除非我将它链接到另一个块或添加消费者?你可以在我发布的原始代码上显示这个吗?
  • 正确。直到项目离开块,它仍被视为“进入”块并计入BoundedCapacity。如果您只想排出项目,可以将转换块连接到缓冲区块:var buffer = new BufferBlock&lt;bool&gt;(); processor.LinkTo(buffer);。或者,如果您确实不需要转换块的输出,请考虑改用 ActionBlock。
猜你喜欢
  • 2011-04-27
  • 1970-01-01
  • 1970-01-01
  • 2015-04-09
  • 1970-01-01
  • 1970-01-01
  • 2023-03-27
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多