【问题标题】:Make sure ProcessingQueue.Count correct in multiple threading application确保 ProcessingQueue.Count 在多线程应用程序中正确
【发布时间】:2023-03-27 21:06:01
【问题描述】:

我有一个 Windows 服务来处理链表队列中的 xml 文件。队列中的文件是在文件创建时由 FileSystemWatcher 事件添加的。

namespace XMLFTP
{
    public class XML_Processor : ServiceBase
    {
       public string s_folder { get; set; }
       public XML_Processor(string folder)
       {
           s_folder = folder;
       }
       Thread worker;
       FileSystemWatcher watcher;
       DirectoryInfo my_Folder;
       public static AutoResetEvent ResetEvent { get; set; }
       bool running;
       public bool Start()
       {
           my_Folder = new DirectoryInfo(s_folder);
           bool success = true;
           running = true;
           worker = new Thread(new ThreadStart(ServiceLoop));
           worker.Start();
           // add files to queue by FileSystemWatcher event
           return (success);
       }
       public bool Stop()
       {
           try
           {
               running = false;
               watcher.EnableRaisingEvents = false;
               worker.Join(ServiceSettings.ThreadJoinTimeOut);
           }
           catch (Exception ex)
           {
               return (false);
           } 
           return (true);
       }
       public void ServiceLoop()
       {
           string fileName;
           while (running)
           {
               Thread.Sleep(2000);
               if (ProcessingQueue.Count > 0)
               {
                   // process file and write info to DB. 
               }
           }
       }

       void watcher_Created(object sender, FileSystemEventArgs e)
       {
           switch (e.ChangeType)
           {
               case WatcherChangeTypes.Created:// add files to queue
           }
       } 
    }
 }

可能存在线程安全问题。

        while (running)
        {
            Thread.Sleep(2000);
            if (ProcessingQueue.Count > 0)
            {
                // process file and write info to DB. 
            }
        }

由于对 ProcessingQueue.Count 的访问不受锁保护,因此如果不同的线程更改“队列”,则 Count 可能会更改。结果,进程文件部分可能会失败。如果您将 Count 属性实现为:

public static int Count
{
get { lock (syncRoot) return _files.Count; }
}

因为锁被提前释放。

我的两个问题:

  1. 如何使ProcessingQueue.Count正确?
  2. 如果我使用.NET Framework 4.5 BlockingCollection技能,示例代码为:

     class ConsumingEnumerableDemo
     {
        // Demonstrates: 
        //      BlockingCollection<T>.Add() 
        //      BlockingCollection<T>.CompleteAdding() 
        //      BlockingCollection<T>.GetConsumingEnumerable() 
        public static void BC_GetConsumingEnumerable()
        {
            using (BlockingCollection<int> bc = new BlockingCollection<int>())
            {
    
                // Kick off a producer task
                Task.Factory.StartNew(() =>
                {
                    for (int i = 0; i < 10; i++)
                    {
                        bc.Add(i);
                        Thread.Sleep(100); // sleep 100 ms between adds
                    }
    
                    // Need to do this to keep foreach below from hanging
                    bc.CompleteAdding();
                });
    
                // Now consume the blocking collection with foreach. 
                // Use bc.GetConsumingEnumerable() instead of just bc because the 
                // former will block waiting for completion and the latter will 
                // simply take a snapshot of the current state of the underlying collection. 
                foreach (var item in bc.GetConsumingEnumerable())
                {
                    Console.WriteLine(item);
                }
             }
           }
         }
    

示例使用常量 10 作为迭代子句,如何将队列中的动态计数应用于它?

【问题讨论】:

    标签: c# multithreading


    【解决方案1】:

    使用BlockingCollection,您不必知道计数。消费者知道继续处理项目,直到队列为空并且IsCompleted 为真。所以你可以有这个:

    var producer = Task.Factory.StartNew(() =>
    {
        // Add 10 items to the queue
        foreach (var i in Enumerable.Range(0, 10))
            queue.Add(i);
    
        // Wait one minute
        Thread.Sleep(TimeSpan.FromMinutes(1.0));
    
        // Add 10 more items to the queue
        foreach (var i in Enumerable.Range(10, 10))
            queue.Add(i);
    
        // mark the queue as complete for adding
        queue.CompleteAdding();
    });
    
    // consumer
    foreach (var item in queue.GetConsumingEnumerable())
    {
        Console.WriteLine(item);
    }
    

    消费者将输出前 10 个项目,这会清空队列。但是因为生产者还没有调用CompleteAdding,所以消费者会继续阻塞队列。它将捕获生产者写入的下 10 个项目。然后,队列为空且IsCompleted == true,因此消费者结束(GetConsumingEnumerable 到达队列的末尾)。

    您可以随时查看Count,但您获得的值只是一个快照。当您评估它时,生产者或消费者很可能已经修改了队列并更改了计数。但这应该没关系。只要你不调用CompleteAdding,消费者就会继续等待物品。

    生产者写入的项目数不必是恒定的。例如,在我的Simple Multithreading 博客文章中,我展示了一个生产者读取文件并将项目写入由消费者提供服务的BlockingCollection。生产者和消费者同时运行,一切都进行到生产者到达文件末尾。

    【讨论】:

    • @Love: GetConsumingEnumerable 检查IsCompleted。这是BlockingCollection 类的属性。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-01-14
    • 2015-09-29
    • 2012-03-26
    • 2011-12-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多