【问题标题】:Process TCP Socket data read in different thread处理在不同线程中读取的 TCP Socket 数据
【发布时间】:2020-11-11 06:32:12
【问题描述】:

我正在开发一个程序来与每 70 毫秒通过 TCP 通道发送数据的设备集成。

我正在使用 Socket.BeginReceive 和 Socket.EndReceive 方法来读取数据。逻辑描述在下面的伪代码中

private void OnReceived(IAsyncResult ar)
{
    var rcvdDataLength = m_tcpSocket.EndReceive(ar);
    Array.Copy(m_tempRecvBuffer, 0, m_mainBuffer, m_mainBufferDataIndex, rcvdDataLength);
    if (CheckIfValidHeaderAndBodyReceived())
    {
        var actualData = new byte[headerLen + BodyLen];
        Array.Copy(m_mainBuffer, m_dataIndex, actualData, 0, headerLen + BodyLen);
        Process(actualData);
    }
    m_tcpSocket.BeginReceive(m_tempRecvBuffer, 0, m_tempRecvBuffer.Length, SocketFlags.None,
        OnReceived, null);
}

上面代码中描述的Process函数负责实现业务逻辑。处理函数目前大约需要 300 毫秒。

因此消费者(需要 300 毫秒)比生产者(每 70 毫秒发布一次数据)慢。我是否需要异步运行此 Process 函数以避免延迟?还是 TCP 层的flow control 方面负责这个?

【问题讨论】:

  • 如果消费时间比接收消息的速率长约 3 倍,您似乎将不得不考虑丢弃一些消息。简单地将处理移动到不同的线程会导致线程数量不断增加,直到资源耗尽,除非我错过了什么。
  • 如果您可以并行化消费者(请参阅BlockingCollection<T>),以便您能够以足够快的速度消费数据(例如,5 个消费线程同时工作,平均可以处理传入的数据每个数据单元的净速率为 60 毫秒),这可以工作。另一方面,如果消费端点直到当前数据单元被处理后才开始新的读取操作,那么是的...... TCP 协议将在本地“处理”事情。但是,远程端点最终可能会耗尽其资源,具体取决于它自己的实现。
  • 底线:可能有很多不同的结果,也有很多不同的解决方案可能是合适的,而且您的帖子中没有足够的信息让任何人都可以提供一个好的答案。
  • 不要混淆从套接字读取数据和在其他地方处理数据。
  • 就像我说的,这取决于远程端点是如何实现的。但是,如果它每 70 毫秒生成一次数据并无条件地尝试将该数据传输到您的客户端,并且客户端仅每 300 毫秒接收一次数据单元,那么正在传输的数据将在服务器排队并最终消耗服务器的内存。一些服务器实现会检测无响应/太慢客户端的情况,并且会丢弃数据或完全丢弃客户端。其他人可能不会。

标签: c# .net sockets tcp winsock


【解决方案1】:

我是否需要异步运行这个 Process 函数以避免延迟?

这取决于您的申请,最终取决于您自己的优先级决定。严格来说,不,您不需要异步执行任何操作,但这可能是件好事。

我通常使用和推荐的方法是一个专用于接口的线程,它通过队列与应用程序的其余部分进行交互。当通信线程接收到消息时,它会锁定一个队列并将它们推入。当主应用程序准备好使用该数据时,它会锁定该队列并尽可能多地退出队列。这是一个简单、健壮且可靠的机制。转到您的伪代码:

private void OnReceived(IAsyncResult ar)
{
    var rcvdDataLength = m_tcpSocket.EndReceive(ar);
    Array.Copy(m_tempRecvBuffer, 0, m_mainBuffer, m_mainBufferDataIndex, rcvdDataLength);
    if (CheckIfValidHeaderAndBodyReceived())
    {
        var actualData = new byte[headerLen + BodyLen];
        Array.Copy(m_mainBuffer, m_dataIndex, actualData, 0, headerLen + BodyLen);
        lock(messageQueue)
        {
           messageQueue.Enqueue(actualData);
        }
    }
    m_tcpSocket.BeginReceive(m_tempRecvBuffer, 0, m_tempRecvBuffer.Length, SocketFlags.None,
        OnReceived, null);
}

然后在你的应用程序的某个地方:

void ProcessQueue()
{
   Queue<byte[]> tempQueue = new Queue<byte[]>();
   lock(messageQueue)
   {
      // Drain the queue so we can release the lock ASAP
      while(messageQueue.Count > 0)
      {
         tempQueue.Enqueue(messageQueue.Dequeue());
      }
   }
   while(tempQueue.Count > 0)
   {
      Process(tempQueue.Dequeue());
   }
}

【讨论】:

  • 感谢您的回答。我会试试这个并回复你
  • 您的解决方案非常有效。你能建议如何运行 ProcessQueue 吗?我实际上是在它自己的专用线程中的无限循环中调用它。或者我应该在一个时间间隔很短(比如 30 毫秒)的计时器中调用 ProcessQueue?
  • 这是另一个架构问题,取决于应用程序的目标和现有结构。您提到消费者需要 300 毫秒,所以我假设消费者在准备好时会调用 ProcessQueue。
  • FWIW,我本周第一次将这种方法用于 TCP。我传统上使用这种队列方法进行串行和 USB 通信。我有一个 BackgroundWorker 管理 TcpListener,接收和解码消息,并将有效消息放入队列。另一个 BackgroundWorker 监视队列并在将消息放入不同的 IEnumerable 之前对消息执行一些缓慢的加密处理。当然会有开销,但它可以确保像 RSA 解密这样的缓慢过程不会阻塞网络接口或 GUI。
猜你喜欢
  • 1970-01-01
  • 2012-03-22
  • 2017-04-27
  • 2011-08-09
  • 2021-12-03
  • 2012-05-24
  • 1970-01-01
  • 1970-01-01
  • 2012-12-02
相关资源
最近更新 更多