【问题标题】:How to copy stream to many streams async C#.NET如何将流复制到异步 C#.NET 的多个流
【发布时间】:2011-11-08 15:31:24
【问题描述】:

我有 TCP 服务器,可以不间断地接收大数据。 而且我需要将此流广播给许多客户。

更新: 我需要播放视频流。也许有现成的解决方案?

【问题讨论】:

  • 请提供有关您尝试过的代码以及哪些代码有效/无效的更多详细信息。
  • 广播有必要使用TCP吗?您的要求通常是通过 UDP 实现的...
  • 我尝试广播的数据 - 它是流。
  • 如果我可以从 silverlight 应用程序连接到 UPD,则不需要 TCP。

标签: c# .net asynchronous stream


【解决方案1】:

如果您想异步执行此操作,则可以利用System.Threading.Tasks namespace

首先,您需要将Stream 实例映射到可以等待完成的Task

IDictionary<Stream, Task> streamToTaskMap = outputStreams.
    ToDictionary(s => s, Task.Factory.StartNew(() => { });

上面有一点开销,因为有一个浪费的 Task 实例什么都不做,但考虑到您需要执行的 Task 实例和延续的数量,这个代价很小。

从那里,您将从流中读取内容,然后将其异步写入每个 Stream 实例:

byte[] buffer = new byte[<buffer size>];
int read = 0;

while ((read = inputStream.Read(buffer, 0, buffer.Length)) > 0)
{
    // The buffer to copy into.
    byte[] copy = new byte[read];

    // Perform the copy.
    Array.Copy(buffer, copy, read);

    // Cycle through the map, and replace the task with a continuation
    // on the task.
    foreach (Stream stream in streamToTaskMap.Keys)
    {
        // Continue.
        streaToTaskMap[stream] = streaToTaskMap[stream].ContinueWith(t => {
            // Write the bytes from the copy.
            stream.Write(copy, 0, copy.Length);
        });
    }
}

最后,您可以通过调用等待所有写入的流:

Task.WaitAll(streamToTaskMap.Values.ToArray());

有几点需要注意。

首先,由于传递给ContinueWith 的lambda,所以需要buffer 的副本; lambda 是一个封装buffer 的闭包,因为它是异步处理的,所以内容可能会发生变化。每个延续都需要自己的缓冲区副本才能读取。

这也是对Stream.Write 的调用使用Array.Length 属性的原因;否则,read 变量必须通过循环的每次迭代复制。

另外,在Stream 类上使用BeginWrite/EndWrite 方法会更理想;因为没有 ContinueWithAsync 方法会采用 Task 并继续使用异步方法,所以调用 read 的异步版本没有任何好处。

在这种情况下,最好自己调用 BeginWrite/EndWrite(以及 BeginRead/EndRead)以充分利用异步操作;当然,这会更复杂一些,因为您不会封装Task 提供的操作结果,并且如果您使用匿名方法/闭包,则必须对buffer 采取相同的预防措施。

【讨论】:

    【解决方案2】:

    产生一个将流传递给它的线程,以及要写入的流

    byte[] buffer = new byte[BUFFER_SIZE];
    int btsRead = 0;
    
    while ((btsRead = inputStream.Read(buffer, 0, BUFFER_SIZE)) > 0)
    {
        foreach (Stream oStream in outputStreams)
            oStream.Write(buffer, 0, btsRead);
    }
    

    编辑:并行写入:

    将 foreach 块替换为:

    Parallel.ForEach(outputStreams, oStream =>
    {
        oStream.Write(buffer, 0, btsRead);
    });
    

    【讨论】:

    • “foreach”?是很好的解决方案吗?如果我有大约 1000 个客户怎么办?
    • foreach 与 for 循环相同。我刚刚展示的代码应该在服务器接收到上传流时在后台线程中运行。您可以尝试其他技术,例如对客户端流进行异步写入以进一步分摊成本,但我提出的是一个相对简单的解决方案,应该可以工作。
    【解决方案3】:

    基本上,您希望线程化您的应用程序。这是一个简单的线程和 TCP/IP 示例

    C# Tutorial - Simple Threaded TCP Server | Switch on the Code

    【讨论】:

      猜你喜欢
      • 2017-06-10
      • 2015-10-27
      • 2010-09-18
      • 2016-05-22
      • 2019-09-14
      • 1970-01-01
      • 1970-01-01
      • 2021-10-01
      相关资源
      最近更新 更多