基于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
相关产品推荐
相关产品推荐

