【问题标题】:Queue calls to asynchronous method BeginSend of a socket对套接字的异步方法 BeginSend 的调用进行排队
【发布时间】:2014-05-04 12:02:27
【问题描述】:

我需要将 BeginSend 调用排队到一个套接字,并且我需要它们按时间顺序执行。为此,我使用了一个信号量来指示回调函数何时可以执行。
大多数情况下,它之所以有效,是因为每个异步回调都在一个单独的线程上执行,但偶尔在新的异步调用中使用当前回调中使用的同一线程。发生这种情况时,该线程被锁定,等待信号量被释放,但是因为应该清除信号量的同一个线程正在等待它被清除,所以该线程被永远锁定。

为了说明问题,这里有一段测试代码:

static Semaphore semaphore = new Semaphore(1, 1);
static IList<byte[]> buffer = new List<byte[]>();
static void Main()
{
    Socket socket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
    socket.Connect(new IPEndPoint(new IPAddress(new byte[] { 192, 168, 1, 8 }), 123));
    while (true) // data feed
    {
        lock (buffer)
        {
            buffer.Add(new byte[1460]);
            if (buffer.Count == 1)
                socket.BeginSend(buffer[0], 0, 1460, 0, new AsyncCallback(SendCallback), socket); // calls BeginSend if the buffer was empty before 
        }
    }
}

static void SendCallback(IAsyncResult ar)
{
    Console.WriteLine("in " + Thread.CurrentThread.ManagedThreadId);
    semaphore.WaitOne();
    Socket socket = (Socket)ar.AsyncState;
    lock(buffer)
    {
        buffer.RemoveAt(0); // removes data that was sent
        if (buffer.Count > 0) // if there is more data to send calls BeginSend again
            socket.BeginSend(buffer[0], 0, 1460, 0, new AsyncCallback(SendCallback), socket);
    }
    semaphore.Release();
    Console.WriteLine("out " + Thread.CurrentThread.ManagedThreadId);
}

这是输出:

因为线程 10 被转移到一个新的回调,没有给前一个回调机会退出和清除信号量,线程被永远锁定。

我应该如何解决这个问题?

【问题讨论】:

  • 不知道 Socket.BeginRead 可以同步调用回调。这会产生可怕的重入问题。你确定吗?您可以发布调用堆栈的屏幕截图来证明这一点吗?在 SendCallback 中应该有两个堆栈帧。
  • 使用 BlockingCollection 来跟踪要发送的数据。解决2个问题:维护秩序,解除僵局
  • 确实如此。这是一个可怕的陷阱。可能是以效率的名义添加的。这也使得 BeginSend 不是非阻塞的,因为任意代码可以作为调用的一部分运行。
  • 你现在有一个竞争条件,因为你正在访问锁之外的缓冲区。

标签: c# multithreading sockets asynchronous


【解决方案1】:

切换到任务:

不错的msdn文章Tasks and the APM Pattern

public Task<int> SendAsync(Socket socket, byte[] buffer, int offset, int size, SocketFlags flags)
{
   var result = socket.BeginSend(buffer, offset, size, flags, _ => { }, socket);
   return Task.Factory.FromAsync(result,(r) => socket.EndSend(r));
}

现在事情变得简单了一点:

使用默认的BlockingCollection<> 作为并发队列。它是线程保存并删除列表上的显式锁定

static BlockingCollection<byte[]> buffer = new BlockingCollection<byte[]>();

public async void Main()
{
   Socket socket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
   socket.Connect(new IPEndPoint(new IPAddress(new byte[] { 192, 168, 1, 8 }), 123));
   while (!buffer.IsCompleted)
   {
      var data = buffer.Take();
      await SendAsync(socket, data, 0, data.Length, 0);
   }
   Console.ReadLine();            
}

在不需要信号量的同时保持非阻塞发送和顺序。

【讨论】:

  • 我不知道任务会改变什么。任务完成仍然可以是同步的并导致重入问题。您已将问题结构更改为避免问题的 while 循环,但可能是 OP 需要锁定并防止并发写入的原因。
  • @usr 由 OP 更新问题并解释原因,也许我们可以调整答案以更好地满足他的需求。
  • 在我真正的问题中,一个线程在它们到达缓冲区列表时逐个添加数据包,如果之前缓冲区中没有数据包,则调用 BeginSend。然后如果有更多的数据包要发送,BeginSend 回调函数会再次调用 BeginSend。不确定任务将如何提供帮助。它基本上会以同步方法转换 BeginSend。
  • @Chris,它看起来像一个同步方法,但实际上不是。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2014-06-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多