【问题标题】:How do I correctly read and write with a Socket/NetworkStream using C# async/await in a single thread?如何在单线程中使用 C# async/await 正确读取和写入 Socket/NetworkStream?
【发布时间】:2018-04-09 15:20:45
【问题描述】:

我正在尝试用 C# 为专有 TCP 协议连接(发送/接收键+值对消息)编写客户端。我想使用NetworkStream 的async/await 功能,所以我的程序可以使用单线程(a la JavaScript)读取和写入套接字但是我有NetworkStream.ReadAsync 工作方式的问题。

这是我的大纲:

public static async Task Main(String[] args)
{
    using( TcpClient tcp = new TcpClient() )
    {
        await tcp.ConnectAsync( ... );
        using( NetworkStream ns = tcp.GetStream() )
        {
            while( true )
            {
                await RunInnerAsync( ns );
            }
        }
    }
}

private static readonly ConcurrentQueue<NameValueCollection> pendingMessagesToSend = new ConcurrentQueue<NameValueCollection>();

private static async Task RunInnerAsync( NetworkStream ns )
{
    // 1. Send any pending messages.
    // 2. Read any newly received messages.

    // 1:
    while( !pendingMessagesToSend.IsEmpty )
    {
        // ( foreach Message, send down the NetworkStream here )
    }

    // 2:
    Byte[] buffer = new Byte[1024];
    while( ns.DataAvailable )
    {
        Int32 bytesRead = await ns.ReadAsync( buffer, 0, buffer.Length );
        if( bytesRead == 0 ) break;
        // ( process contents of `buffer` here )
    } 
}

这里有一个问题:如果NetworkStream ns 中没有要读取的数据(DataAvailable == false),那么Main 中的while( true ) 循环会不断运行,CPU 永远不会空闲 - 这很糟糕。

因此,如果我更改代码以删除DataAvailable 检查并始终调用ReadAsync,则调用有效地“阻塞”直到数据可用 - 因此,如果没有数据到达,则此客户端将永远不会向远程主机。于是想着加个500ms左右的超时时间:

// 2:
Byte[] buffer = new Byte[1024];
while( ns.DataAvailable )
{
    CancellationTokenSource cts = new CancellationTokenSource( 500 );
    Task<Int32> readTask = ns.ReadAsync( buffer, 0, buffer.Length, cts.Token );
    await readTask;
    if( readTask.IsCancelled ) break;
    // ( process contents of `buffer` here )
}

However, this does not work!显然,当实际请求取消时,接受 CancellationToken 的 NetworkStream.ReadAsync 重载不会中止或停止,它总是忽略它(这不是错误吗?)。

我链接到的 QA 建议了一个简单地关闭 Socket/NetworkStream 的解决方法——这对我来说是不合适的,因为我需要保持连接处于活动状态,但只有 休息一下等待数据到达并发送一些数据。

One of the other answers 建议共同等待Task.Delay,如下所示:

// 2:
Byte[] buffer = new Byte[1024];
while( ns.DataAvailable )
{
    Task maxReadTime = Task.Delay( 500 );
    Task readTask = ns.ReadAsync( buffer, 0, buffer.Length );
    await Task.WhenAny( maxReadTime, readTask );
    if( maxReadTime.IsCompleted )
    {
        // what do I do here to cancel the still-pending ReadAsync operation?
    }
}

...然而,虽然这确实阻止了程序等待无限期的网络读取操作,它本身并没有停止读取操作 - 所以当我的程序完成发送任何未决消息时,它会在它仍在等待数据到达时再次调用ReadAsync - 这意味着要处理重叠 IO,这根本不是我想要的。

I know when working with Socket directly and its BeginReceive / EndReceive methods 你只需要在EndReceive 中调用BeginReceive - 但是第一次如何安全地调用BeginReceive,尤其是在循环中 - 以及在使用@ 时应该如何修改这些调用987654349@ API 代替?

【问题讨论】:

  • 也许你应该添加SpinWait 结构而不是while (true),所以最终主线程会产生执行?

标签: sockets asynchronous async-await networkstream


【解决方案1】:

我通过从NetworkStream.ReadAsync 获取Task 并将其存储在可变局部变量中来使其工作.

不幸的是,async 方法不能有ref 参数,我需要将我的逻辑从RunInnerAsync 移动到Main。

这是我的解决方案:

private static readonly ConcurrentQueue<NameValueCollection> pendingMessagesToSend = new ConcurrentQueue<NameValueCollection>();
private static TaskCompletionSource<Object> queueTcs;

public static async Task Main(String[] args)
{
    using( TcpClient tcp = new TcpClient() )
    {
        await tcp.ConnectAsync( ... );
        using( NetworkStream ns = tcp.GetStream() )
        {
            Task<Int32> nsReadTask = null; // <-- this!
            Byte[] buffer = new Byte[1024];

            while( true )
            {
                if( nsReadTask == null )
                {
                    nsReadTask = ns.ReadAsync( buffer, 0, buffer.Length );
                }

                if( queueTcs == null ) queueTcs = new TaskCompletionSource<Object>();

                Task completedTask = await Task.WhenAny( nsReadTask, queueTcs.Task );
                if( completedTask == nsReadTask )
                {
                    while( ns.DataAvailable )
                    {
                        Int32 bytesRead = await ns.ReadAsync( buffer, 0, buffer.Length );
                        if( bytesRead == 0 ) break;
                        // ( process contents of `buffer` here )
                    } 
                }
                else if( completedTask == queueTcs )
                {
                    while( !pendingMessagesToSend.IsEmpty )
                    {
                        // ( foreach Message, send down the NetworkStream here )
                    }
                }
            }
        }
    }
}

每当pendingMessagesToSend 被修改时,queueTcs 被实例化(如果为 null)并调用 SetResult(null) 以取消等待 Task completedTask = await Task.WhenAny( nsReadTask, queueTcs.Task ); 行。

自从构建了这个解决方案后,我觉得我使用 TaskCompletionSource 并不恰当,请参阅此 QA:Is it acceptable to use TaskCompletionSource as a WaitHandle substitute?

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-08-23
    • 2021-04-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-16
    相关资源
    最近更新 更多