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

Azure Relay混合连接:同步请求-响应实现及请求响应匹配验证

确保Hybrid Connections请求与响应匹配的方案

在使用Azure Relay Hybrid Connections的流模式时,因为没有内置的请求-响应绑定机制,咱们得自己实现关联逻辑,最可靠的方式是给每个请求分配唯一标识,再配合并发安全的请求管理来匹配响应。下面结合你的代码示例给出具体实现思路:

核心思路:用Request ID关联请求与响应

每个请求生成一个唯一的ID(比如Guid),发送请求时把ID和请求内容一起传递给服务端;服务端处理完成后,将这个ID和响应内容一起返回;客户端收到响应后,通过ID找到对应的请求,完成匹配。

具体实现代码

客户端修改后的SendAsync方法

我们会用ConcurrentDictionary来管理待处理的请求(处理并发场景),结合TaskCompletionSource来异步等待对应响应:

using System.Text.Json;
using System.Collections.Concurrent;

private HybridConnectionClient _client = new HybridConnectionClient("你的Hybrid Connection连接字符串");
// 存储待处理请求:Key是RequestId,Value是等待响应的TaskCompletionSource
private readonly ConcurrentDictionary<Guid, TaskCompletionSource<RelayResponse>> _pendingRequests = new ConcurrentDictionary<Guid, TaskCompletionSource<RelayResponse>>();

public override async Task<RelayResponse> SendAsync(RelayRequest request)
{
    // 生成唯一请求ID
    var requestId = Guid.NewGuid();
    var responseTcs = new TaskCompletionSource<RelayResponse>();
    _pendingRequests.TryAdd(requestId, responseTcs);

    try
    {
        var stream = await _client.CreateConnectionAsync();
        var writer = new StreamWriter(stream) { AutoFlush = true };
        var reader = new StreamReader(stream);

        // 把RequestId和请求内容打包成可序列化的对象
        var wrappedRequest = new WrappedRelayRequest
        {
            RequestId = requestId,
            Request = request
        };
        var requestJson = JsonSerializer.Serialize(wrappedRequest);
        
        // 发送请求(用WriteLine确保每个请求是独立的一行,避免粘包)
        await writer.WriteLineAsync(requestJson);

        // 读取服务端返回的响应
        var responseJson = await reader.ReadLineAsync();
        var wrappedResponse = JsonSerializer.Deserialize<WrappedRelayResponse>(responseJson);

        // 验证响应的RequestId是否匹配当前请求
        if (wrappedResponse?.RequestId == requestId && _pendingRequests.TryRemove(requestId, out var tcs))
        {
            tcs.SetResult(wrappedResponse.Response);
            return wrappedResponse.Response;
        }
        else
        {
            var error = new InvalidOperationException("收到不匹配的响应,无法关联到当前请求");
            responseTcs.SetException(error);
            throw error;
        }
    }
    catch (Exception ex)
    {
        // 出现异常时移除待处理请求,避免内存泄漏
        _pendingRequests.TryRemove(requestId, out _);
        responseTcs.SetException(ex);
        throw;
    }
}

// 辅助包装类,用于携带RequestId
public class WrappedRelayRequest
{
    public Guid RequestId { get; set; }
    public RelayRequest Request { get; set; }
}

public class WrappedRelayResponse
{
    public Guid RequestId { get; set; }
    public RelayResponse Response { get; set; }
}

关键补充:超时处理

为了避免请求无限等待,建议给每个请求添加超时逻辑,修改上面的代码,在发送请求后加入超时判断:

// 发送请求后,等待响应或超时
var timeoutTask = Task.Delay(TimeSpan.FromSeconds(30)); // 30秒超时
var completedTask = await Task.WhenAny(responseTcs.Task, timeoutTask);

if (completedTask == timeoutTask)
{
    _pendingRequests.TryRemove(requestId, out _);
    var timeoutEx = new TimeoutException("请求超时,未收到响应");
    responseTcs.SetException(timeoutEx);
    throw timeoutEx;
}

return await responseTcs.Task;

服务端配合逻辑

服务端必须正确识别RequestId,并在返回响应时携带它,示例代码如下:

public async Task HandleConnectionAsync(Stream stream)
{
    var reader = new StreamReader(stream);
    var writer = new StreamWriter(stream) { AutoFlush = true };

    while (true)
    {
        try
        {
            // 读取客户端发送的请求
            var requestJson = await reader.ReadLineAsync();
            if (string.IsNullOrEmpty(requestJson)) break;

            var wrappedRequest = JsonSerializer.Deserialize<WrappedRelayRequest>(requestJson);
            // 处理业务请求
            var relayResponse = await ProcessRelayRequest(wrappedRequest.Request);

            // 包装响应并返回(携带原RequestId)
            var wrappedResponse = new WrappedRelayResponse
            {
                RequestId = wrappedRequest.RequestId,
                Response = relayResponse
            };
            var responseJson = JsonSerializer.Serialize(wrappedResponse);
            await writer.WriteLineAsync(responseJson);
        }
        catch (Exception ex)
        {
            // 处理异常,可返回带RequestId的错误响应
            Console.WriteLine($"处理请求出错:{ex.Message}");
            break;
        }
    }
}

// 模拟业务处理方法
private Task<RelayResponse> ProcessRelayRequest(RelayRequest request)
{
    // 你的业务逻辑
    return Task.FromResult(new RelayResponse { /* 填充响应内容 */ });
}

额外注意事项

  • 避免粘包:示例中用ReadLineAsync和WriteLineAsync来分隔每个请求/响应,确保数据边界清晰;如果用二进制数据,可以用固定长度前缀来标识数据长度。
  • 连接复用:如果要复用Hybrid Connection连接(而不是每次请求创建新连接),需要确保在同一个流上的请求/响应能正确按顺序关联,Request ID机制依然适用。
  • 并发安全:ConcurrentDictionary是线程安全的,适合处理多并发请求的场景,避免请求混乱。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 08:07:02