我发现拥有一个处理协调和写入然后使用异步处理读取的单个线程是最简单的。通过在AsyncState 对象中包含CancellationTokenSource,当接收方遇到错误(调用EndRead 或流完成)时,阅读器代码可以向发送线程发出信号以重新启动连接。协调/编写器位于创建连接的循环中,然后循环使用BlockingCollection<T> 的要发送的项目。通过使用BlockingCollection.GetConsumingEnumerable(token),可以在阅读器遇到错误时取消发送者。
private class AsyncState
{
public byte[] Buffer { get; set; }
public NetworkStream NetworkStream { get; set; }
public CancellationTokenSource CancellationTokenSource { get; set; }
}
一旦您创建了连接,您就可以开始异步读取过程(只要一切正常,它就会一直调用自己)。在状态对象中传递缓冲区、流和CancellationTokenSource:
var buffer = new byte[1];
stream.BeginRead(buffer, 0, 1, Callback,
new AsyncState
{
Buffer = buffer,
NetworkStream = stream,
CancellationTokenSource = cts2
});
之后,您开始从输出队列中读取数据并写入流,直到取消或发生故障:
using (var writer = new StreamWriter(stream, Encoding.ASCII, 80, true))
{
foreach (var item in this.sendQueue.GetConsumingEnumerable(cancellationToken))
{
...
...在回调中,您可以检查故障,如有必要,点击 CancellationTokenSource 以向编写器线程发出信号以重新启动连接。
private void Callback(IAsyncResult ar)
{
var state = (AsyncState)ar.AsyncState;
if (ar.IsCompleted)
{
try
{
int bytesRead = state.NetworkStream.EndRead(ar);
LogState("Post read ", state.NetworkStream);
}
catch (Exception ex)
{
Log.Warn("Exception during EndRead", ex);
state.CancellationTokenSource.Cancel();
return;
}
// Deal with the character received
char c = (char)state.Buffer[0];
if (c < 0)
{
Log.Warn("c < 0, stream closing");
state.CancellationTokenSource.Cancel();
return;
}
... deal with the character here, building up a buffer and
... handing it out to the application when completed
... perhaps using Rx Subject<T> to make it easy to subscribe
... and finally ask for the next byte with the same Callback
// Launch the next reader
var buffer2 = new byte[1];
var state2 = state.WithNewBuffer(buffer2);
state.NetworkStream.BeginRead(buffer2, 0, 1, Callback, state2);