【问题标题】:Best practices for implementing a thread to do fast, bulk, and continuous reading in C#?在 C# 中实现线程以进行快速、批量和连续读取的最佳实践?
【发布时间】:2015-02-06 08:42:44
【问题描述】:

在 .NET 4.0 中应该如何处理从 C# 中的设备读取批量数据?具体来说,我需要从一个 USB HID 设备中快速读取,该设备发出超过 26 个数据包的报告,其中必须保留顺序。

我尝试在 BackgroundWorker 线程中执行此操作。它一次从设备读取一个数据包,并在读取更多数据之前对其进行处理。这提供了相当好的响应时间,但它可能会在这里和那里丢失一个数据包,并且单个数据包读取的开销成本会增加。

while (!( sender as BackgroundWorker ).CancellationPending) {
       //read a single packet
       //check for header or footer
       //process packet data
    }
}

在 C# 中读取这样的设备的最佳做法是什么?


背景:

我的 USB HID 设备不断报告大量数据。数据分成 26 个数据包,我必须保留订单。不幸的是,该设备只标记每个报告中的第一个最后一个数据包,因此我需要能够捕获其间的所有其他数据包。

【问题讨论】:

  • 什么版本的.net?答案将取决于它。
  • @MatthewWatson 我的目标是 .NET 4.0,但如果一个答案也能解释与其他版本的区别,那就太好了。
  • HID 数据速率非常低,最高 8 KB/秒。你不能编写跟不上它的代码,不需要“最佳实践”。
  • @HansPassant 这比每毫秒 64 字节的情况要高一点,但这个问题与 USB 或 HID 无关。我只是好奇应该如何处理这样的事情。

标签: c# multithreading io usb hid


【解决方案1】:

对于 .Net 4,您可以使用 BlockingCollection 来提供可供生产者和消费者使用的线程安全队列。 BlockingCollection.GetConsumingEnumerable() 方法提供了一个枚举器,当队列使用CompleteAdding() 标记为已完成且为空时,该枚举器自动终止。

这里有一些示例代码。在此示例中,有效负载是一个整数数组,但当然您可以使用您需要的任何数据类型。

请注意,对于您的具体示例,您可以使用the overload of GetConsumingEnumerable(),它接受CancellationToken 类型的参数。

using System;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;

namespace Demo
{
    public static class Program
    {
        private static void Main()
        {
            var queue = new BlockingCollection<int[]>();

            Task.Factory.StartNew(() => produce(queue));

            consume(queue);

            Console.WriteLine("Finished.");
        }

        private static void consume(BlockingCollection<int[]> queue)
        {
            foreach (var item in queue.GetConsumingEnumerable())
            {
                Console.WriteLine("Consuming " + item[0]);
                Thread.Sleep(25);
            }
        }

        private static void produce(BlockingCollection<int[]> queue)
        {
            for (int i = 0; i < 1000; ++i)
            {
                Console.WriteLine("Producing " + i);
                var payload = new int[100];
                payload[0] = i;
                queue.Add(payload);
                Thread.Sleep(20);
            }

            queue.CompleteAdding();
        }
    }
}

对于 .Net 4.5 及更高版本,您可以使用来自 Microsoft's Task Parallel Library 的更高级别的类,它具有丰富的功能(乍一看可能有些令人生畏)。

这是使用 TPL DataFlow 的相同示例:

using System;
using System.Threading;
using System.Threading.Tasks;
using System.Threading.Tasks.Dataflow;

namespace Demo
{
    public static class Program
    {
        private static void Main()
        {
            var queue = new BufferBlock<int[]>();

            Task.Factory.StartNew(() => produce(queue));
            consume(queue).Wait();

            Console.WriteLine("Finished.");
        }

        private static async Task consume(BufferBlock<int[]> queue)
        {
            while (await queue.OutputAvailableAsync())
            {
                var payload = await queue.ReceiveAsync();
                Console.WriteLine("Consuming " + payload[0]);
                await Task.Delay(25);
            }
        }

        private static void produce(BufferBlock<int[]> queue)
        {
            for (int i = 0; i < 1000; ++i)
            {
                Console.WriteLine("Producing " + i);
                var payload = new int[100];
                payload[0] = i;
                queue.Post(payload);
                Thread.Sleep(20);
            }

            queue.Complete();
        }
    }
}

【讨论】:

  • 我正要发布一个非常相似的答案,唯一不同的是我让我的 produceconsume 接受了 CancellationToken,这样我就可以重新创建 OP 的行为(!( sender as BackgroundWorker ).CancellationPending)
  • 我做了足够多的工作,我还是决定发布我的答案。
  • 很好的答案,感谢您也发布了 4.5 的示例。
【解决方案2】:

如果担心丢失数据包,请不要在同一线程上进行处理和读取。从 .NET 4.0 开始,他们添加了 System.Collections.Concurrent 命名空间,这使得这很容易做到。您所需要的只是一个BlockingCollection,它充当传入数据包的队列。

BlockingCollection<Packet> _queuedPackets = new BlockingCollection<Packet>(new ConcurrentQueue<Packet>());

void readingBackgroundWorker_DoWork(object sender, DoWorkEventArgs e)
{
    while (!( sender as BackgroundWorker ).CancellationPending) 
    {
       Packet packet = GetPacket();
       _queuedPackets.Add(packet);
    }        
    _queuedPackets.CompleteAdding();
}

void processingBackgroundWorker_DoWork(object sender, DoWorkEventArgs e)
{
    List<Packet> report = new List<Packet>();
    foreach(var packet in _queuedPackets.GetConsumingEnumerable())
    {
        report.Add(packet);
        if(packet.IsLastPacket)
        {
            ProcessReport(report);
            report = new List<Packet>();
        }
    }
}

_queuedPackets 为空时会发生_queuedPackets.GetConsumingEnumerable() 将阻塞线程而不消耗任何资源。一旦数据包到达,它将解除阻塞并执行 foreach 的下一次迭代。

当您调用_queuedPackets.CompleteAdding(); 时,处理线程上的 foreach 将运行直到集合为空,然后退出 foreach 循环。如果您不希望它在取消时“完成队列”,您可以轻松地将其更改为提前退出。我还将改用任务而不是后台工作人员,因为它使传递参数更容易。

void ReadingLoop(BlockingCollection<Packet> queue, CancellationToken token)
{
    while (!token.IsCancellationRequested) 
    {
       Packet packet = GetPacket();
       queue.Add(packet);
    }        
    queue.CompleteAdding();
}

void ProcessingLoop(BlockingCollection<Packet> queue, CancellationToken token)
{
    List<Packet> report = new List<Packet>();

    try
    {
        foreach(var packet in queue.GetConsumingEnumerable(token))
        {
            report.Add(packet);
            if(packet.IsLastPacket)
            {
                ProcessReport(report);
                report = new List<Packet>();
            }
        }
    }
    catch(OperationCanceledException)
    {
        //Do nothing, we don't care that it happened.
    }
}

//This would replace your backgroundWorker.RunWorkerAsync() calls;
private void StartUpLoops()
{
    var queue = new BlockingCollection<Packet>(new ConcurrentQueue<Packet>());    
    var cancelRead = new CancellationTokenSource();
    var cancelProcess = new CancellationTokenSource();
    Task.Factory.StartNew(() => ReadingLoop(queue, cancelRead.Token));
    Task.Factory.StartNew(() => ProcessingLoop(queue, cancelProcess.Token));

    //You can stop each loop indpendantly by calling cancelRead.Cancel() or cancelProcess.Cancel()
}

【讨论】:

  • 感谢您决定发布此内容,这是非常有用的 +1
猜你喜欢
  • 2011-04-17
  • 1970-01-01
  • 1970-01-01
  • 2012-01-07
  • 1970-01-01
  • 2010-10-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多