【问题标题】:How to work threading with ConcurrentQueue<T>如何使用 ConcurrentQueue<T> 处理线程
【发布时间】:2011-05-31 20:49:12
【问题描述】:

我正在尝试找出使用队列的最佳方式。我有一个返回 DataTable 的进程。每个 DataTable 依次与前一个 DataTable 合并。有一个问题,在最终 BulkCopy (OutOfMemory) 之前保存的记录太多。

所以,我决定我应该立即处理每个传入的 DataTable。考虑ConcurrentQueue&lt;T&gt;...但我不知道WriteQueuedData() 方法如何知道将表出列并将其写入数据库。

例如:

public class TableTransporter
{
    private ConcurrentQueue<DataTable> tableQueue = new ConcurrentQueue<DataTable>();

    public TableTransporter()
    {
        tableQueue.OnItemQueued += new EventHandler(WriteQueuedData);   // no events available
    }

    public void ExtractData()
    {
        DataTable table;

        // perform data extraction
        tableQueue.Enqueue(table);
    }

    private void WriteQueuedData(object sender, EventArgs e)
    {
        BulkCopy(e.Table);
    }
}

我的第一个问题是,除了我实际上没有要订阅的任何事件这一事实之外,如果我异步调用ExtractData(),这就是我所需要的吗?其次,关于ConcurrentQueue&lt;T&gt; 的运行方式以及需要某种形式的触发器来与排队的对象异步工作,我是否遗漏了什么?

更新 我刚刚从ConcurrentQueue&lt;T&gt; 派生了一个具有 OnItemQueued 事件处理程序的类。那么:

new public void Enqueue (DataTable Table)
{
    base.Enqueue(Table);
    OnTableQueued(new TableQueuedEventArgs(Table));
}

public void OnTableQueued(TableQueuedEventArgs table)
{
    EventHandler<TableQueuedEventArgs> handler = TableQueued;

    if (handler != null)
    {
        handler(this, table);
    }
}

对此实现有任何顾虑吗?

【问题讨论】:

标签: c# multithreading queue concurrent-collections


【解决方案1】:

根据我对问题的理解,您遗漏了一些东西。

并发队列是一种数据结构,旨在接受多个线程读取和写入队列,而无需显式锁定数据结构。 (所有爵士乐都在幕后处理,或者集合以不需要锁定的方式实现。)

考虑到这一点,您尝试使用的模式似乎是“生产者/消费者”。首先,您有一些任务产生工作(并将项目添加到队列中)。其次,您还有第二个任务消耗队列中的东西(并从队列中取出项目)。

所以你真的需要两个线程:一个添加项目,第二个删除项目。因为您使用的是并发集合,所以可以有多个线程添加项目和多个线程删除项目。但显然你在并发队列上的争用越多,就会越快成为瓶颈。

【讨论】:

  • 我以为我有 2 个线程。主线程基本上会等待事件触发。第二个线程以对ExtractData() 的异步调用开始。在异步回调中,我将继续提取过程。
  • 其实我觉得我倒退了;主线程应该排队数据表;然后通过入队的项目事件触发器开始异步写入方法。
  • @Chris Smith:我不知道怎么给你发消息。您的个人资料中有一个恶意网站。请删除它。
  • @KamranBigdely 感谢您指出这一点!我博客的域名失效了,坏演员接管了它。固定。
【解决方案2】:

我认为ConcurrentQueue 仅在极少数情况下有用。它的主要优点是它是无锁的。然而,通常生产者线程必须以某种方式通知消费者线程有数据可供处理。线程之间的这种信号需要锁定,并否定了使用ConcurrentQueue 的好处。同步线程最快的方法是使用Monitor.Pulse(),它只在锁内有效。所有其他同步工具甚至更慢。

当然,消费者可以不断地检查队列中是否有东西,这在没有锁的情况下工作,但是对处理器资源的浪费是巨大的。如果消费者在检查之间等待会更好一点。

在写入队列时引发线程是一个非常糟糕的主意。使用ConcurrentQueue 节省1 微秒可能会被执行eventhandler 完全浪费,这可能需要1000 倍的时间。

如果所有处理都在事件处理程序或异步调用中完成,那么问题是为什么还需要队列?最好将数据直接传递给处理程序,并且根本不使用队列。

请注意ConcurrentQueue 的实现相当复杂以允许并发。在大多数情况下,最好使用普通的Queue&lt;&gt; 并锁定对队列的每次访问。由于队列访问只需要微秒,因此两个线程在同一微秒内访问队列的可能性极小,并且几乎不会因为锁定而出现任何延迟。使用带锁定的普通Queue&lt;&gt; 通常会比ConcurrentQueue 更快地执行代码。

【讨论】:

  • 对收到反对票感到羞耻。我认为这是一个有效的、务实的意见。
  • >生产者线程必须以某种方式通知消费者线程有数据可用于处理您通常如何执行此操作?
  • 有关线程如何相互同步的概述,请参阅:Microsoft, Overview of Synchronization Primitives。在这些情况下,ConcurrentQueue 没有帮助,因为 Synchronization Primitives 无论如何都使用锁定。 ConcurrentQueue 可能对大量并行问题很有用,当几个线程尽可能快地产生某些东西并且另一个线程收集这些结果并处理它们时。问题解决后,释放所有线程,无需等待。
【解决方案3】:

这是我想出的完整解决方案:

public class TableTransporter
{
    private static int _indexer;

    private CustomQueue tableQueue = new CustomQueue();
    private Func<DataTable, String> RunPostProcess;
    private string filename;

    public TableTransporter()
    {
        RunPostProcess = new Func<DataTable, String>(SerializeTable);
        tableQueue.TableQueued += new EventHandler<TableQueuedEventArgs>(tableQueue_TableQueued);
    }

    void tableQueue_TableQueued(object sender, TableQueuedEventArgs e)
    {
        //  do something with table
        //  I can't figure out is how to pass custom object in 3rd parameter
        RunPostProcess.BeginInvoke(e.Table,new AsyncCallback(PostComplete), filename);
    }

    public void ExtractData()
    {
        // perform data extraction
        tableQueue.Enqueue(MakeTable());
        Console.WriteLine("Table count [{0}]", tableQueue.Count);
    }

    private DataTable MakeTable()
    { return new DataTable(String.Format("Table{0}", _indexer++)); }

    private string SerializeTable(DataTable Table)
    {
        string file = Table.TableName + ".xml";

        DataSet dataSet = new DataSet(Table.TableName);

        dataSet.Tables.Add(Table);

        Console.WriteLine("[{0}]Writing {1}", Thread.CurrentThread.ManagedThreadId, file);
        string xmlstream = String.Empty;

        using (MemoryStream memstream = new MemoryStream())
        {
            XmlSerializer xmlSerializer = new XmlSerializer(typeof(DataSet));
            XmlTextWriter xmlWriter = new XmlTextWriter(memstream, Encoding.UTF8);

            xmlSerializer.Serialize(xmlWriter, dataSet);
            xmlstream = UTF8ByteArrayToString(((MemoryStream)xmlWriter.BaseStream).ToArray());

            using (var fileStream = new FileStream(file, FileMode.Create))
                fileStream.Write(StringToUTF8ByteArray(xmlstream), 0, xmlstream.Length + 2);
        }
        filename = file;

        return file;
    }

    private void PostComplete(IAsyncResult iasResult)
    {
        string file = (string)iasResult.AsyncState;
        Console.WriteLine("[{0}]Completed: {1}", Thread.CurrentThread.ManagedThreadId, file);

        RunPostProcess.EndInvoke(iasResult);
    }

    public static String UTF8ByteArrayToString(Byte[] ArrBytes)
    { return new UTF8Encoding().GetString(ArrBytes); }

    public static Byte[] StringToUTF8ByteArray(String XmlString)
    { return new UTF8Encoding().GetBytes(XmlString); }
}

public sealed class CustomQueue : ConcurrentQueue<DataTable>
{
    public event EventHandler<TableQueuedEventArgs> TableQueued;

    public CustomQueue()
    { }
    public CustomQueue(IEnumerable<DataTable> TableCollection)
        : base(TableCollection)
    { }

    new public void Enqueue (DataTable Table)
    {
        base.Enqueue(Table);
        OnTableQueued(new TableQueuedEventArgs(Table));
    }

    public void OnTableQueued(TableQueuedEventArgs table)
    {
        EventHandler<TableQueuedEventArgs> handler = TableQueued;

        if (handler != null)
        {
            handler(this, table);
        }
    }
}

public class TableQueuedEventArgs : EventArgs
{
    #region Fields
    #endregion

    #region Init
    public TableQueuedEventArgs(DataTable Table)
    {this.Table = Table;}
    #endregion

    #region Functions
    #endregion

    #region Properties
    public DataTable Table
    {get;set;}
    #endregion
}

作为概念证明,它似乎运作良好。我最多看到 4 个工作线程。

【讨论】:

  • 看这个,这是一个很好的实现,但是,在运行了一个快速测试之后,一个项目什么时候出队?
  • @RichardPriddy:因为那是 5 年前的事了(我早就搬到了我的第三家公司),我只能假设这不是一个完整的例子.请注意末尾的概念证明备注。 ;) 也就是说,根据要求,您可以公开enqueued 事件并让其他东西处理出列。否则,在后处理函数的AsyncCallback 中的某处出列可能是合乎逻辑的。在这么晚的日子里,真的很难找到更具体的东西。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-02-27
  • 2013-05-25
  • 1970-01-01
  • 2016-08-29
相关资源
最近更新 更多