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

如何在BlockingCollection被填充时立即调用消费者方法?

基于BlockingCollection实现无轮询的即时数据处理

背景

通过查阅资料了解到,BlockingCollection<T>的设计初衷就是消除线程间共享集合需要主动轮询检查新数据的麻烦。当有新数据插入集合时,消费者线程会被立即唤醒,无需在while循环里按固定间隔查询数据。

需求

  • 一个容量为1的BlockingCollection;
  • 3个生产者位置向集合填充数据;
  • 当前用while循环轮询检查集合是否有数据;
  • 希望集合一有数据就立即执行ProcessInbox()并清空集合,不需要固定间隔轮询。

现有代码

using System;
using System.Collections.Concurrent;
using System.Linq;
using System.Threading;
        
namespace ConsoleApp1
{
     class Program
     {
          private static BlockingCollection<int> _processingNotificationQueue = new(1);

          private static void GetDataFromQueue(CancellationToken cancellationToken)
          {
               Console.WriteLine("GDFQ called");
               int data;
               //while (!cancellationToken.IsCancellationRequested)
               while(!_processingNotificationQueue.IsCompleted)
               {
                    try
                    {
                         if(_processingNotificationQueue.TryTake(out data))
                         {
                              Console.WriteLine("Take");
                              ProcessInbox();
                         }
                    }
                    catch (Exception ex)
                    {

                    }
               }
          }
        
          private static void ProcessInbox()
          {
               Console.WriteLine("PI called");
          }
        
          private static void PostDataToQueue(object state)
          {
               Console.WriteLine("PDTQ called");
               _processingNotificationQueue.TryAdd(1);
          }
        
          private void MessageInsertedToTabale()
          {
               PostDataToQueue(new CancellationToken());
          }
        
          private void FewMessagesareNotProcessed()
          {
               PostDataToQueue(new CancellationToken());
          }
        
          static void Main(string[] args)
          {
               Console.WriteLine("Start");
               new Timer(PostDataToQueue, new CancellationToken(), TimeSpan.Zero,
                   TimeSpan.FromMilliseconds(100));
        
               // new Thread(()=> PostDataToQueue()).Start();
               new Thread(() => GetDataFromQueue(new CancellationToken())).Start();
        
               Console.WriteLine("End");
               Console.ReadKey();
          }
     }
}

解决方案

核心是利用BlockingCollection<T>的阻塞式消费机制替代轮询逻辑,让消费者线程在无数据时休眠,有数据时立即被唤醒处理。

修改后的消费者方法

private static void GetDataFromQueue(CancellationToken cancellationToken)
{
    Console.WriteLine("GDFQ called");
    try
    {
        // GetConsumingEnumerable会自动阻塞等待数据,无需轮询
        foreach (var data in _processingNotificationQueue.GetConsumingEnumerable(cancellationToken))
        {
            Console.WriteLine("Take");
            ProcessInbox();
            // 容量为1,每次Take后集合自动为空,无需额外清空操作
        }
    }
    catch (OperationCanceledException)
    {
        Console.WriteLine("Consumer operation canceled");
    }
}

关键改动说明

  1. 替换轮询逻辑:用GetConsumingEnumerable替代while+TryTake的轮询,该方法会在集合为空时自动阻塞线程,直到有数据加入或取消令牌触发,彻底消除轮询开销。
  2. 优雅控制线程退出:结合CancellationToken处理线程终止,比单纯依赖IsCompleted更灵活可靠,避免线程无法正常停止的问题。
  3. 移除无效异常捕获:原代码空的catch块会隐藏错误,保留针对取消操作的异常捕获即可。

生产者端优化(可选)

因为集合容量为1,TryAdd可能失败导致数据丢失,可改用阻塞式Add方法确保数据入队:

private static void PostDataToQueue(object state)
{
    Console.WriteLine("PDTQ called");
    var token = (CancellationToken)state;
    try
    {
        // 阻塞直到集合有空间或被取消
        _processingNotificationQueue.Add(1, token);
    }
    catch (OperationCanceledException)
    {
        Console.WriteLine("Producer operation canceled");
    }
}

主方法调整

添加取消令牌管理,确保程序退出时能正确终止所有线程:

static void Main(string[] args)
{
    Console.WriteLine("Start");
    using var cts = new CancellationTokenSource();
    // 传递取消令牌给定时器
    new Timer(PostDataToQueue, cts.Token, TimeSpan.Zero, TimeSpan.FromMilliseconds(100));
    
    new Thread(() => GetDataFromQueue(cts.Token)).Start();
    
    Console.WriteLine("Press any key to exit...");
    Console.ReadKey();
    // 发送取消信号终止所有操作
    cts.Cancel();
}

内容的提问来源于stack exchange,提问作者Dreamer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 01:05:24