如何让Order Controller在MassTransit Saga执行完所有消费者后收到响应?
问题:请求响应模式下工作流完成后API超时无响应
我采用请求响应模式确保工作流按顺序执行,现有三个项目:Order.API、OrderPayment.API、OrderConfirmation.API。所有消费者能正常执行,但最后一个消费者完成后,Swagger持续加载,每次抛出请求超时错误,Order Controller无法收到执行成功或失败的响应,直至返回请求超时错误。相关代码如下:
OrderPayment.API 消费者代码
public class PaymentConsumer : IConsumer<PaymentProcessed> { public async Task Consume(ConsumeContext<PaymentProcessed> context) { Console.WriteLine($"Payment processed for OrderId: {context.Message.OrderId}"); await context.Send(new OrderConfirmed { OrderId = context.Message.OrderId }); } }
OrderPayment.API Program.cs 代码
builder.Services.AddMassTransit(x => { x.AddConsumer<PaymentConsumer>(); x.UsingRabbitMq((context, cfg) => { cfg.Host("localhost", "/", h => { h.Username("guest"); h.Password("guest");}); cfg.ConfigureEndpoints(context); }); });
OrderConfirmation.API 消费者代码
public class ConfirmationConsumer : IConsumer<OrderConfirmed> { public async Task Consume(ConsumeContext<OrderConfirmed> context) { Console.WriteLine($"Order {context.Message.OrderId} has been confirmed."); await Task.CompletedTask; } }
OrderConfirmation.API Program.cs 代码
builder.Services.AddMassTransit(x => { x.AddConsumer<ConfirmationConsumer>(); x.UsingRabbitMq((context, cfg) => { cfg.Host("localhost", "/", h => { h.Username("guest"); h.Password("guest");}); cfg.ConfigureEndpoints(context); }); });
共享模型、状态及事件代码
namespace SharedModels.OrderProcessModels { // 订单创建事件 public class OrderCreated { public Guid OrderId { get; set; } } // 订单创建响应 public class OrderCreatedResponse { public Guid OrderId { get; set; } public string Status { get; set; } = "OrderCreated"; } // 支付处理完成事件 public class PaymentProcessed { public Guid OrderId { get; set; } } // 支付处理完成响应 public class PaymentProcessedResponse { public Guid OrderId { get; set; } public string Status { get; set; } = "PaymentProcessed"; } // 订单确认事件 public class OrderConfirmed { public Guid OrderId { get; set; } } // 订单确认响应 public class OrderConfirmedResponse { public Guid OrderId { get; set; } public string Status { get; set; } = "OrderConfirmed"; } // 跟踪订单生命周期的状态对象 public class OrderState : SagaStateMachineInstance { public Uri ResponseAddress { get; set; } public Guid RequestId { get; set; } // 可选,存储RequestId用于关联 public Guid CorrelationId { get; set; } // Saga实例唯一标识 public Guid OrderId { get; set; } // 订单标识 public State CurrentState { get; set; } // 订单在Saga中的当前状态 public DateTime? CreatedAt { get; set; } // 订单创建时间 public DateTime? PaymentAt { get; set; } // 支付处理完成时间 public DateTime? ConfirmedAt { get; set; } // 订单确认时间 } }
Order.API OrderStateMachine.cs 代码
public class OrderStateMachine : MassTransitStateMachine<OrderState> { public State PaymentProcessed { get; private set; } public State OrderConfirmed { get; private set; } public Event<OrderCreated> OrderCreatedEvent { get; private set; } public Event<PaymentProcessed> PaymentProcessedEvent { get; private set; } public Event<OrderConfirmed> OrderConfirmedEvent { get; private set; } public OrderStateMachine() { InstanceState(x => x.CurrentState); // 通过OrderId关联事件与Saga实例 Event(() => OrderCreatedEvent, x => x.CorrelateById(c => c.Message.OrderId)); Event(() => PaymentProcessedEvent, x => x.CorrelateById(c => c.Message.OrderId)); Event(() => OrderConfirmedEvent, x => x.CorrelateById(c => c.Message.OrderId)); // 初始状态 -> 支付处理完成状态 Initially( When(OrderCreatedEvent) .Then(context => { context.Saga.OrderId = context.Message.OrderId; // 保存订单ID context.Saga.CreatedAt = DateTime.UtcNow; // 设置创建时间戳 Console.WriteLine("Order created, transitioning to PaymentProcessed..."); }) .TransitionTo(PaymentProcessed) // 切换到PaymentProcessed状态 .Publish(context => new PaymentProcessed { OrderId = context.Saga.OrderId }) ); // 支付处理完成状态 -> 订单确认状态 During(PaymentProcessed, When(PaymentProcessedEvent) .Then(context => { context.Saga.PaymentAt = DateTime.UtcNow; // 设置支付时间戳 Console.WriteLine("Payment processed, transitioning to OrderConfirmed..."); }) .TransitionTo(OrderConfirmed) // 切换到OrderConfirmed状态 .Publish(context => new OrderConfirmed { OrderId = context.Saga.OrderId }) ); // 最终状态 -> 完成Saga During(OrderConfirmed, When(OrderConfirmedEvent) .Then(context => { context.Saga.ConfirmedAt = DateTime.UtcNow; // 设置确认时间戳 Console.WriteLine($"Order confirmed for OrderId: {context.Saga.OrderId}"); }) .Finalize() // 订单确认后终结Saga ); SetCompletedWhenFinalized(); } }
Order.API OrderController.cs 代码
namespace Order.API.Controllers { [ApiController] [Route("api")] public class OrderController : ControllerBase { private readonly IRequestClient<OrderCreated> _requestClient; // 等待第一个响应后继续执行 public OrderController(IRequestClient<OrderCreated> requestClient) { _requestClient = requestClient; } [HttpPost("PlaceOrder")] public async Task<IActionResult> PlaceOrder([FromBody] CreateOrderModel model) { try { var orderId = Guid.Parse(model.OrderId); var status = await _requestClient.GetResponse<OrderConfirmed>(new OrderCreated { OrderId = orderId }); if (status != null) { return Ok(new { OrderId = orderId, Status = "Order confirmed" }); } else { return StatusCode(500, "Order processing failed."); } } catch (Exception e) { return StatusCode(504, "Timeout waiting for the order confirmation."); } } } public class CreateOrderModel { public string OrderId { get; set; } public string CustomerName { get; set; } public double Amount { get; set; } } }
问题原因分析
- 请求响应不匹配:控制器使用
IRequestClient<OrderCreated>.GetResponse<OrderConfirmed>等待响应,但整个工作流中没有任何环节向请求的响应地址发送OrderConfirmed消息,IRequestClient会一直等待对应请求ID的响应直至超时。 - 事件与响应混淆:
OrderConfirmed是事件而非响应消息,控制器应等待OrderConfirmedResponse而非事件本身,且Saga未处理响应回传逻辑。 - Saga未关联请求上下文:初始请求的
ResponseAddress和RequestId未被Saga保存,导致流程结束时无法向控制器发送响应。
解决方案
1. 修正Saga,保存请求上下文并发送最终响应
修改OrderStateMachine的初始状态和最终状态逻辑:
// 初始状态中保存请求上下文 Initially( When(OrderCreatedEvent) .Then(context => { context.Saga.OrderId = context.Message.OrderId; context.Saga.CreatedAt = DateTime.UtcNow; // 保存请求的响应地址和请求ID context.Saga.ResponseAddress = context.Request.ResponseAddress; context.Saga.RequestId = context.Request.RequestId; Console.WriteLine("Order created, transitioning to PaymentProcessed..."); }) .TransitionTo(PaymentProcessed) .Publish(context => new PaymentProcessed { OrderId = context.Saga.OrderId }) ); // 最终状态发送响应给控制器 During(OrderConfirmed, When(OrderConfirmedEvent) .Then(context => { context.Saga.ConfirmedAt = DateTime.UtcNow; Console.WriteLine($"Order confirmed for OrderId: {context.Saga.OrderId}"); }) .Send(context => context.Saga.ResponseAddress, context => new OrderConfirmedResponse { OrderId = context.Saga.OrderId }) .Finalize() );
2. 修正控制器,接收正确的响应类型
[HttpPost("PlaceOrder")] public async Task<IActionResult> PlaceOrder([FromBody] CreateOrderModel model) { try { var orderId = Guid.Parse(model.OrderId); var response = await _requestClient.GetResponse<OrderConfirmedResponse>(new OrderCreated { OrderId = orderId }); return Ok(new { OrderId = orderId, Status = response.Message.Status }); } catch (RequestTimeoutException) { return StatusCode(504, "Timeout waiting for the order confirmation."); } catch (Exception e) { return StatusCode(500, $"Order processing failed: {e.Message}"); } }
3. 优化事件发送逻辑(可选)
将PaymentConsumer中的context.Send改为Publish,保持事件驱动的一致性:
public async Task Consume(ConsumeContext<PaymentProcessed> context) { Console.WriteLine($"Payment processed for OrderId: {context.Message.OrderId}"); await context.Publish(new OrderConfirmed { OrderId = context.Message.OrderId }); }
4. 配置Saga状态持久化(关键)
在Order.API的Program.cs中添加Saga持久化配置(测试用内存存储,生产环境需替换为数据库):
builder.Services.AddMassTransit(x => { x.AddSagaStateMachine<OrderStateMachine, OrderState>() .InMemoryRepository(); x.UsingRabbitMq((context, cfg) => { cfg.Host("localhost", "/", h => { h.Username("guest"); h.Password("guest"); }); cfg.ConfigureEndpoints(context); }); });
内容的提问来源于stack exchange,提问作者Mikh Dany
相关产品推荐
相关产品推荐

