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

修复C# SSL异步客户端BeginWrite写入冲突问题及优化方案

问题修复与优化方案

一、当前APM模式的错误修复

你遇到的The BeginWrite method cannot be called when another write is pending错误,核心原因是同一时间对同一个SslStream发起了多个未完成的异步写入操作。APM模式下,BeginWrite必须等待上一次写入的EndWrite调用完成后,才能发起下一次写入。

修复步骤:

  1. 维护一个线程安全的发送队列,存放待发送的字节数组
  2. 用锁和标记位控制发送流程,确保同一时间只有一个写入操作在执行
  3. 在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}字节数据");
}

方案优势:

  1. 异步代码结构清晰,避免APM模式下的回调嵌套
  2. SemaphoreSlim自动保证发送操作的串行化,不会出现写入冲突
  3. Channel是.Net中高效的异步队列,适合生产者-消费者场景(比如业务线程生产待发送数据,发送循环异步消费)
  4. 天然支持Protobuf序列化/反序列化,只需将byte[]替换为Protobuf对象的序列化结果即可

内容的提问来源于stack exchange,提问作者user12652992

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 17:45:29