【问题标题】:Collecting Data from multiple IO-based Task from a single Thread-based task从单个基于线程的任务中收集来自多个基于 IO 的任务的数据
【发布时间】:2013-01-22 03:10:33
【问题描述】:

使用 TPL,我如何从多个 IO 源(“无线程”任务)收集结果,并在它们从各自的源进来时将它们合并成一个序列,而不为每个源生成一个基于线程的任务来监控它们?从一个线程轮询源是否安全?

while (true)
{
    try
    {
        IEnumerable<UdpClient> readyChannels = 
            from channel in channels
            where channel.Available > 0
            select channel;

        foreach( UdpClient channel in readyChannels)
        {
           var result = await channel.ReceiveAsync();
           //do something with result like post to dataflow block.
        }
    }
    catch (Exception e)
    {
        throw (e);
    }
    ...

这样的事情怎么样?

【问题讨论】:

  • 这里真的需要while循环吗?只是好奇
  • @Cuong Le - 我想不断地轮询和阅读我的 udp 来源。
  • 可以用Task.WhenAny 或类似的结构来做这种事情,但我认为更好的匹配是为每个UdpClient 公开IObservableISourceBlock。 Rx 或 TPL 数据流更适合像这样的“推送”事件。
  • @Stephen C - IObservable/ISourceBlock 火花线程。我需要从一个线程收集所有数据,而不是多个线程监控不同的来源,因为有很多来源。
  • 上面的代码在它自己的线程中执行,但是每个源都应该从那个线程监控,而不是创建一个新的线程/每个源。我正在尝试保存线程。

标签: c# io task-parallel-library tpl-dataflow


【解决方案1】:

我在这里看到了几个选项:

如果您想启动对ReceiveAsync() 的调用,请将它们设置为对结果执行某些操作(例如,如您所说,发送到数据流块)然后忘记它们,您可以使用ContinueWith()

foreach (var channel in readyChannels)
{
   channel.ReceiveAsync().ContinueWith(task => 
   {
       var result = task.Result;
       //do something with result like post to dataflow block.
   }
}

这样做的一个缺点是您需要在每个延续中处理异常。

可能更好的方法是使用 Stephen Cleary 的 AsyncEx 中的 OrderByCompletion()。这样,您可以一次开始所有读取并在它们完成时对其进行处理:

var tasks = readyChannels.Select(c => c.ReceiveAsync()).OrderByCompletion();

foreach (var task in tasks)
{
   var result = await task;
   //do something with result like post to dataflow block.
}

如果您想限制并行性,另一个有用的选项是使用TransformBlock

var receiveBlock = new TransformBlock<UdpClient, UdpReceiveResult>(
    c => c.ReceiveAsync(),
    new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = degreeOfParallelism });
foreach (var channel in readyChannels)
    receiveBlock.Post(channel);
receiveBlock.Complete();

// set up processing here

await receiveBlock.Completion;

如果你想将结果发送到另一个块,那么上面评论中提到的处理包括简单地将它们链接在一起:

receiveBlock.LinkTo(anotherBlock);

在上述所有情况下,从来没有线程阻塞来监视任何东西。但是调用ReceiveAsync()然后处理结果的代码必须在某处执行。

【讨论】:

  • 有没有一种简单的方法可以删除 UdpClients 或阻止它们接收?
猜你喜欢
  • 1970-01-01
  • 2018-02-27
  • 2023-03-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多