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

基于Azure Service Bus实现MVC应用向微服务带参请求数据

基于Azure Service Bus实现MVC与微服务的请求-响应模式(无HTTP)

核心思路

利用Azure Service Bus原生的**关联ID(CorrelationId)或会话(Session)**机制绑定请求与响应,同时通过精准的消息过滤避免锁定整个队列,实现异步的请求-响应交互,替代HTTP模式。

具体实现方案

1. 定义消息结构

给请求消息添加唯一标识的CorrelationId,同时指定响应的目标队列:

// 请求消息
public class DataRequest
{
    public Guid CorrelationId { get; set; }
    public string ReplyQueue { get; set; }
    // 业务参数
    public int QueryParam { get; set; }
}

// 响应消息
public class DataResponse
{
    public Guid CorrelationId { get; set; }
    public string Result { get; set; }
}

2. MVC端发送请求并监听响应

  • 生成唯一CorrelationId,创建临时回复队列(或复用专用响应队列),发送请求时绑定该ID与回复队列地址;
  • 仅监听带有当前请求CorrelationId的响应消息,避免无关消息干扰:
var correlationId = Guid.NewGuid();
var tempReplyQueue = $"reply-queue-{correlationId:N}";

// 初始化请求消息
var requestMsg = new ServiceBusMessage(JsonSerializer.Serialize(new DataRequest
{
    CorrelationId = correlationId,
    ReplyQueue = tempReplyQueue,
    QueryParam = 456
}))
{
    CorrelationId = correlationId.ToString(),
    ReplyTo = tempReplyQueue
};

// 发送请求到微服务队列
await serviceBusSender.SendMessageAsync(requestMsg);

// 创建响应接收器,仅过滤当前CorrelationId的消息
var receiverOptions = new ServiceBusReceiverOptions
{
    Filter = new SqlFilter($"CorrelationId = '{correlationId.ToString()}'"),
    ReceiveMode = ServiceBusReceiveMode.PeekLock
};

using var replyReceiver = serviceBusClient.CreateReceiver(tempReplyQueue, receiverOptions);
// 等待响应(设置超时)
var responseMsg = await replyReceiver.ReceiveMessageAsync(TimeSpan.FromSeconds(20));

if (responseMsg != null)
{
    var response = JsonSerializer.Deserialize<DataResponse>(responseMsg.Body.ToString());
    // 处理业务结果
    Console.WriteLine(response.Result);
    // 完成消息,释放锁定
    await replyReceiver.CompleteMessageAsync(responseMsg);
}

// 删除临时回复队列
await serviceBusClient.DeleteQueueAsync(tempReplyQueue);

3. 微服务端处理请求并返回响应

微服务监听请求队列,解析请求参数后,将响应发送到指定的回复队列,并保持CorrelationId一致:

await foreach (var requestMsg in requestReceiver.ReceiveMessagesAsync())
{
    try
    {
        var request = JsonSerializer.Deserialize<DataRequest>(requestMsg.Body.ToString());
        // 执行业务逻辑生成结果
        var response = new DataResponse
        {
            CorrelationId = request.CorrelationId,
            Result = $"处理完成:参数{request.QueryParam}"
        };

        // 创建响应消息,绑定CorrelationId
        var responseBusMsg = new ServiceBusMessage(JsonSerializer.Serialize(response))
        {
            CorrelationId = request.CorrelationId.ToString()
        };

        // 发送到回复队列
        var replySender = serviceBusClient.CreateSender(request.ReplyQueue);
        await replySender.SendMessageAsync(responseBusMsg);

        // 完成请求消息,释放锁定
        await requestReceiver.CompleteMessageAsync(requestMsg);
    }
    catch (Exception)
    {
        // 异常处理:放弃消息或重试
        await requestReceiver.AbandonMessageAsync(requestMsg);
    }
}

4. 避免全队列锁定的关键

  • 精准消息过滤:通过SqlFilter或会话ID,让接收器只处理与当前请求绑定的消息,不会锁定队列内其他消息;
  • 临时回复队列:每个请求使用独立的临时队列,避免响应消息混排,进一步缩小锁定范围;
  • PeekLock模式:仅锁定当前处理的单条消息,处理完成后立即释放,不会占用整个队列的锁定资源。

简化方案:使用SDK内置请求-响应API

Azure.Messaging.ServiceBus SDK已封装请求-响应逻辑,自动处理CorrelationId和回复队列,无需手动管理:

// MVC端发送请求
var requestMsg = new ServiceBusMessage(JsonSerializer.Serialize(new { QueryParam = 789 }));
var responseMsg = await serviceBusSender.SendRequestAsync(requestMsg, TimeSpan.FromSeconds(20));
var response = JsonSerializer.Deserialize<DataResponse>(responseMsg.Body.ToString());

// 微服务端处理请求
await foreach (var request in requestReceiver.ReceiveMessagesAsync())
{
    var response = new ServiceBusMessage(JsonSerializer.Serialize(new { Result = "处理完成" }));
    await requestReceiver.CompleteMessageAsync(request);
    await serviceBusSender.SendMessageAsync(response, request.ReplyTo, request.CorrelationId);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 21:05:24