如何实现C# async TCP读写同步?
如何实现C# async TCP读写同步?
嘿,我刚好做过类似的实时TCP通信组件,你的需求其实很典型——用双队列分别管理 inbound(待处理的收到消息)和 outbound(待发送的消息),再结合异步IO和线程安全机制就能完美搞定。下面给你一步步拆解实现思路和核心代码:
核心设计思路
要满足你的需求,核心要解决三个问题:
- 线程安全的消息队列:业务线程、读线程、写线程会同时操作队列,必须保证队列操作的原子性;
- 异步读写分离:TCP的读和写要分开异步执行,避免互相阻塞,保证实时性;
- 有序的消息发送:多个线程可能同时提交发送请求,要保证消息按顺序发送,且不会出现写操作冲突。
具体实现步骤
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
相关产品推荐
相关产品推荐

