修复C# SSL异步客户端BeginWrite写入冲突问题及优化方案
问题修复与优化方案
一、当前APM模式的错误修复
你遇到的The BeginWrite method cannot be called when another write is pending错误,核心原因是同一时间对同一个SslStream发起了多个未完成的异步写入操作。APM模式下,BeginWrite必须等待上一次写入的EndWrite调用完成后,才能发起下一次写入。
修复步骤:
- 维护一个线程安全的发送队列,存放待发送的字节数组
- 用锁和标记位控制发送流程,确保同一时间只有一个写入操作在执行
- 在
EndWrite回调中,检查队列是否还有待发送数据,若有则继续发起下一次BeginWrite
示例修改后的核心代码:
private readonly ConcurrentQueue<byte[]> _sendQueue = new ConcurrentQueue<byte[]>(); private readonly object _sendLock = new object(); private bool _isWriting; public void EnqueueSend(byte[] data) { _sendQueue.Enqueue(data); // 触发发送,若当前没有写入操作则启动 lock (_sendLock) { if (!_isWriting) { _isWriting = true; ProcessSendQueue(); } } } private void ProcessSendQueue() { if (_sendQueue.TryDequeue(out var data)) { _sslStream.BeginWrite(data, 0, data.Length, OnWriteComplete, null); } else { lock (_sendLock) { _isWriting = false; } } } private void OnWriteComplete(IAsyncResult ar) { try { _sslStream.EndWrite(ar); } catch (Exception ex) { // 处理写入异常,比如断开连接、记录日志 Console.WriteLine($"写入失败: {ex.Message}"); } // 继续处理队列中的下一个发送任务 ProcessSendQueue(); } // 原SendWorker方法改为调用EnqueueSend private void SendWorker() { while (_connected) { // 模拟获取待发送数据(比如从业务队列取) if (TryGetDataToSend(out var data)) { EnqueueSend(data); } Thread.Sleep(100); // 根据实际业务调整,避免空循环占用CPU } }
二、更完善的TAP模式方案
.Net Framework 4.7.2完全支持TAP(基于任务的异步模式),相比APM更简洁易维护,且原生支持异步等待,无需手动管理回调和队列。
核心思路:
- 用
SslStream.ReadAsync循环监听服务器回复 - 用
SemaphoreSlim实现发送操作的串行化(避免同时发起多个WriteAsync) - 用
Channel维护发送队列,异步消费待发送数据
示例代码:
private readonly SslStream _sslStream; private readonly SemaphoreSlim _sendSemaphore = new SemaphoreSlim(1, 1); private readonly Channel<byte[]> _sendChannel = Channel.CreateUnbounded<byte[]>(); private bool _isRunning; public async Task StartClientAsync() { _isRunning = true; // 启动接收循环 _ = ReceiveLoopAsync(); // 启动发送消费循环 _ = SendLoopAsync(); } private async Task ReceiveLoopAsync() { var buffer = new byte[4096]; try { while (_isRunning) { var bytesRead = await _sslStream.ReadAsync(buffer, 0, buffer.Length); if (bytesRead == 0) { // 服务器断开连接 _isRunning = false; break; } // 处理接收到的数据(比如反序列化为Protobuf对象) ProcessReceivedData(buffer, bytesRead); } } catch (Exception ex) { Console.WriteLine($"接收异常: {ex.Message}"); _isRunning = false; } } private async Task SendLoopAsync() { try { await foreach (var data in _sendChannel.Reader.ReadAllAsync()) { await _sendSemaphore.WaitAsync(); try { await _sslStream.WriteAsync(data, 0, data.Length); await _sslStream.FlushAsync(); } finally { _sendSemaphore.Release(); } } } catch (Exception ex) { Console.WriteLine($"发送异常: {ex.Message}"); } } // 对外提供的发送方法 public void SendData(byte[] data) { if (_isRunning) { _sendChannel.Writer.TryWrite(data); } } private void ProcessReceivedData(byte[] buffer, int bytesRead) { // 这里处理服务器回复,比如用Protobuf反序列化 // var message = Serializer.Deserialize<YourProtobufType>(new MemoryStream(buffer, 0, bytesRead)); Console.WriteLine($"收到{bytesRead}字节数据"); }
方案优势:
- 异步代码结构清晰,避免APM模式下的回调嵌套
SemaphoreSlim自动保证发送操作的串行化,不会出现写入冲突Channel是.Net中高效的异步队列,适合生产者-消费者场景(比如业务线程生产待发送数据,发送循环异步消费)- 天然支持Protobuf序列化/反序列化,只需将byte[]替换为Protobuf对象的序列化结果即可
内容的提问来源于stack exchange,提问作者user12652992
相关产品推荐
相关产品推荐

