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

如何为Thrift netstd .NET Standard添加WebSocket传输支持

为C# .NET Standard客户端库添加WebSocket传输支持的实现方案

你提到的继承TStreamTransport并基于System.Net.WebSockets实现的思路是正确的,但还有几个关键要点需要注意:

  • WebSocket连接生命周期管理

    • 必须手动处理连接建立(握手超时、CancellationToken控制)、异常断开重连、主动关闭流程。比如调用ClientWebSocket.ConnectAsync时要设置合理超时,断开后要先调用CloseAsync发送标准关闭帧,再释放资源,避免直接Dispose导致连接异常。
    • 要让TStreamTransport的OpenAsync/CloseAsync方法与WebSocket的状态流转对应,确保连接状态一致。
  • 流与WebSocket消息的适配

    • TStreamTransport依赖流的读写模型,而ClientWebSocket是基于消息的,需要封装一个WebSocketStream类继承Stream,将WebSocket的消息收发转换为流的读写操作:
      • 读取时,将ReceiveAsync获取的消息缓存到内存流,供ReadAsync分段读取;
      • 写入时,将流数据打包成WebSocket二进制/文本消息发送,注意设置EndOfMessage标记。
  • .NET Standard兼容性处理

    • System.Net.WebSockets.ClientWebSocket仅在.NET Standard 2.0及以上版本支持,需确保库的目标框架符合要求。
    • 不同平台(如.NET Framework、.NET Core)对WebSocket的实现存在差异,比如部分平台不支持某些扩展协议,需做兼容性判断。
  • 错误处理与日志

    • 捕获WebSocket传输中的异常(网络中断、握手失败、超时等),按照客户端库的错误机制抛出或处理,避免未捕获异常导致库崩溃。
    • 添加关键节点日志(连接建立/断开、消息收发量),方便问题排查。
  • 性能优化

    • 若业务场景允许,实现连接池复用WebSocket连接,避免频繁创建销毁连接的开销。
    • 处理大消息分片,设置合理的接收缓存大小,防止内存溢出。

简单实现示例

public class WebSocketStream : Stream
{
    private readonly ClientWebSocket _webSocket;
    private readonly MemoryStream _receiveBuffer = new MemoryStream();
    private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(1, 1);

    public WebSocketStream(ClientWebSocket webSocket)
    {
        _webSocket = webSocket ?? throw new ArgumentNullException(nameof(webSocket));
    }

    public override async Task<int> ReadAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
    {
        await _semaphore.WaitAsync(cancellationToken);
        try
        {
            if (_receiveBuffer.Position >= _receiveBuffer.Length)
            {
                _receiveBuffer.SetLength(0);
                var tempBuffer = new byte[4096];
                WebSocketReceiveResult result;
                do
                {
                    result = await _webSocket.ReceiveAsync(new ArraySegment<byte>(tempBuffer), cancellationToken);
                    await _receiveBuffer.WriteAsync(tempBuffer, 0, result.Count, cancellationToken);
                } while (!result.EndOfMessage);

                _receiveBuffer.Position = 0;
            }

            return await _receiveBuffer.ReadAsync(buffer, offset, count, cancellationToken);
        }
        finally
        {
            _semaphore.Release();
        }
    }

    public override async Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
    {
        await _webSocket.SendAsync(
            new ArraySegment<byte>(buffer, offset, count), 
            WebSocketMessageType.Binary, 
            true, 
            cancellationToken);
    }

    // 实现其他Stream抽象方法
    public override void Flush() => throw new NotSupportedException();
    public override int Read(byte[] buffer, int offset, int count) => throw new NotSupportedException();
    public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException();
    public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException();
    public override void SetLength(long value) => throw new NotSupportedException();
    public override bool CanRead => true;
    public override bool CanSeek => false;
    public override bool CanWrite => true;
    public override long Length => throw new NotSupportedException();
    public override long Position { get => throw new NotSupportedException(); set => throw new NotSupportedException(); }
}

public class WebSocketTransport : TStreamTransport
{
    private readonly Uri _serverUri;
    private ClientWebSocket _webSocket;
    private WebSocketStream _stream;

    public WebSocketTransport(Uri serverUri)
    {
        _serverUri = serverUri ?? throw new ArgumentNullException(nameof(serverUri));
        _webSocket = new ClientWebSocket();
        // 可添加自定义请求头、子协议等配置
        // _webSocket.Options.SetRequestHeader("X-Client-Version", "1.0");
    }

    protected override async Task<Stream> OpenStreamAsync(CancellationToken cancellationToken)
    {
        if (_webSocket.State != WebSocketState.Open)
        {
            await _webSocket.ConnectAsync(_serverUri, cancellationToken);
        }
        _stream = new WebSocketStream(_webSocket);
        return _stream;
    }

    protected override async Task CloseStreamAsync(Stream stream, CancellationToken cancellationToken)
    {
        if (_webSocket.State == WebSocketState.Open)
        {
            await _webSocket.CloseAsync(
                WebSocketCloseStatus.NormalClosure, 
                "Closed by transport", 
                cancellationToken);
        }
        _webSocket.Dispose();
        _stream.Dispose();
        _webSocket = new ClientWebSocket();
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 21:10:31