如果您想异步执行此操作,则可以利用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 采取相同的预防措施。