You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现C# async TCP读写同步?

如何实现C# async TCP读写同步?

嘿,我刚好做过类似的实时TCP通信组件,你的需求其实很典型——用双队列分别管理 inbound(待处理的收到消息)和 outbound(待发送的消息),再结合异步IO和线程安全机制就能完美搞定。下面给你一步步拆解实现思路和核心代码:

核心设计思路

要满足你的需求,核心要解决三个问题:

  1. 线程安全的消息队列:业务线程、读线程、写线程会同时操作队列,必须保证队列操作的原子性;
  2. 异步读写分离:TCP的读和写要分开异步执行,避免互相阻塞,保证实时性;
  3. 有序的消息发送:多个线程可能同时提交发送请求,要保证消息按顺序发送,且不会出现写操作冲突。

具体实现步骤

1. 选择线程安全的队列

直接用.NET内置的ConcurrentQueue<T>就行,它已经封装了线程安全的入队/出队操作,不用我们手动加锁,省心又高效。

2. 封装异步TCP读写循环

把读和写拆成两个独立的异步循环,分别处理 inbound 和 outbound 消息:

  • 读循环:持续从Socket流读取数据,解析成消息后加入 inbound 队列;
  • 写循环:等待 outbound 队列有消息,批量取出并发送到Socket,发送完成后继续等待新消息。

3. 用信号量控制写循环触发

为了避免写循环空转浪费资源,我们用SemaphoreSlim做信号量:当业务线程把消息加入 outbound 队列时,释放一个信号量,写循环收到信号后就开始处理消息,处理完再等待下一个信号。

核心代码示例

using System;
using System.Collections.Concurrent;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using System.Threading.Tasks;

public class AsyncTcpClient
{
    private readonly TcpClient _tcpClient;
    // 待发送消息队列
    private readonly ConcurrentQueue<byte[]> _outboundQueue = new ConcurrentQueue<byte[]>();
    // 待处理消息队列
    private readonly ConcurrentQueue<string> _inboundQueue = new ConcurrentQueue<string>();
    // 控制写循环的信号量
    private readonly SemaphoreSlim _writeSemaphore = new SemaphoreSlim(0, int.MaxValue);
    private bool _isRunning;
    private CancellationTokenSource _cts;

    public AsyncTcpClient(string host, int port)
    {
        _tcpClient = new TcpClient(host, port);
    }

    // 启动TCP通信循环
    public async Task StartAsync()
    {
        _isRunning = true;
        _cts = new CancellationTokenSource();
        // 同时启动读、写两个异步循环,互不阻塞
        var readTask = RunReadLoopAsync(_cts.Token);
        var writeTask = RunWriteLoopAsync(_cts.Token);
        await Task.WhenAll(readTask, writeTask);
    }

    // 优雅停止通信
    public async Task StopAsync()
    {
        _isRunning = false;
        _cts.Cancel();
        _tcpClient.Close();
        await Task.CompletedTask;
    }

    // 业务层调用:发送消息
    public void SendMessage(string message)
    {
        var data = Encoding.UTF8.GetBytes(message);
        _outboundQueue.Enqueue(data);
        // 释放信号量,通知写循环有消息要发
        _writeSemaphore.Release();
    }

    // 业务层调用:获取收到的消息
    public bool TryGetInboundMessage(out string message)
    {
        return _inboundQueue.TryDequeue(out message);
    }

    // 异步读循环:持续读取Socket数据
    private async Task RunReadLoopAsync(CancellationToken token)
    {
        var stream = _tcpClient.GetStream();
        var buffer = new byte[4096]; // 根据实际消息大小调整缓冲区
        while (_isRunning && !token.IsCancellationRequested)
        {
            try
            {
                var bytesRead = await stream.ReadAsync(buffer, 0, buffer.Length, token);
                if (bytesRead == 0)
                {
                    // 对方关闭连接
                    _isRunning = false;
                    break;
                }
                // 这里假设消息是UTF8字符串,实际项目要根据协议解析(比如长度前缀、固定格式)
                var message = Encoding.UTF8.GetString(buffer, 0, bytesRead);
                _inboundQueue.Enqueue(message);
            }
            catch (IOException)
            {
                // 处理连接异常,比如断开重连、记录日志
                _isRunning = false;
                break;
            }
        }
    }

    // 异步写循环:批量发送队列中的消息
    private async Task RunWriteLoopAsync(CancellationToken token)
    {
        var stream = _tcpClient.GetStream();
        while (_isRunning && !token.IsCancellationRequested)
        {
            // 等待信号量,没有消息时会阻塞在这里,不浪费CPU
            await _writeSemaphore.WaitAsync(token);
            
            // 批量处理所有待发送消息
            while (_outboundQueue.TryDequeue(out var data))
            {
                try
                {
                    await stream.WriteAsync(data, 0, data.Length, token);
                    await stream.FlushAsync(token);
                }
                catch (IOException)
                {
                    // 发送失败可以把消息重新入队或者记录错误
                    _outboundQueue.Enqueue(data);
                    _isRunning = false;
                    break;
                }
            }
        }
    }
}

关键细节说明

  • 粘包处理:上面的示例没有处理TCP粘包问题,如果你的消息是自定义协议(比如带长度前缀),要在RunReadLoopAsync里先读取消息长度,再读取对应长度的字节,避免把多个消息当成一个处理;
  • 错误重试:实际项目中,发送失败或者连接断开时,可以加入重连逻辑,把未发送的消息重新入队;
  • 性能优化:如果消息量极大,可以考虑用Channel<T>代替ConcurrentQueue,Channel提供了更高效的异步生产者/消费者模式,支持背压控制;
  • 线程安全:ConcurrentQueue已经保证了队列操作的线程安全,而SemaphoreSlim则保证了写循环不会被重复触发,同时避免了写操作的冲突。

备注:内容来源于stack exchange,提问作者MarengoHue

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.13 16:48:12