【问题标题】:C# How to coordinate read and write threads while reconnecting NetworkStream?C#如何在重新连接NetworkStream时协调读写线程?
【发布时间】:2015-09-07 04:41:41
【问题描述】:

我有一个线程向 NetworkStream 写入请求。另一个线程正在从该流中读取响应。

我想让它容错。万一出现网络故障,我希望将 NetworkStream 替换为全新的。

我让两个线程处理 IO/Socket 异常。他们每个人都会尝试重新建立连接。我正在努力协调这两个线程。我不得不放置锁定部分,使代码相当复杂且容易出错。

有没有推荐的方法来实现这个?也许使用单线程但使读取或写入异步?

【问题讨论】:

    标签: c# multithreading networkstream


    【解决方案1】:

    我发现拥有一个处理协调和写入然后使用异步处理读取的单个线程是最简单的。通过在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);
    

    【讨论】:

    • 这就是我想探索的路线。不过,我没有为我的用例描绘代码。我需要读取包含正文长度的两个字节标题,然后读取正文。我怎么做?我虽然有 ReadHeader 用 ReadBody 调用 BeingRead 但听起来我需要在 ReadBody 方法中同步读取正文的其余部分。另外,我应该如何使用 CancellationTokenSource?
    • 我添加了代码来展示如何处理一次读取一个字节以及如何传入取消令牌并在发生读取失败时向写入器线程发出信号。您可以先读取两个字节进行初始异步读取,然后在一次调用中读取剩余部分。
    • 太棒了。我真的需要state2吗?为什么我不能重用状态?
    • 是的,可能可以重用状态——出于习惯,我倾向于使用不可变的数据结构。
    猜你喜欢
    • 1970-01-01
    • 2019-12-01
    • 1970-01-01
    • 2020-01-12
    • 2020-12-24
    • 2011-10-31
    • 1970-01-01
    • 1970-01-01
    • 2021-06-24
    相关资源
    最近更新 更多