NServiceBus Saga使用While循环实现延迟投递的问题及修复建议
问题根源分析
- 无限死循环:代码中
while (!Data.IsAllPaymentsDone)没有在循环内更新Data.IsAllPaymentsDone的逻辑,导致线程持续执行循环体,疯狂生成并发送消息,直接引发内存占用过高、应用卡顿。 - 变量引用错误:创建的支付请求对象是
paymentRequest,但发送时传入的是paymentOrchestrationRequest,变量不匹配导致支付处理器接收不到正确的消息。 - 不符合业务的延迟逻辑:需求是每日持续检查待缴费用,但当前用20分钟延迟+循环的方式,既无法实现每日检查的目标,还会产生大量无效的延迟消息,浪费资源。
修复方案
1. 修正变量引用错误
将发送时的参数改为正确的paymentRequest:
await context.Send(paymentRequest, options);
2. 替换无限循环为NServiceBus Saga的延迟调度机制
Saga本身支持通过延迟消息实现定时触发,不需要手动写while循环。正确的做法是:
- 每次检查后,若仍有待缴费用,调度一个延迟24小时的消息触发下一次检查;
- 检查逻辑放在单独的消息处理方法中,避免阻塞线程。
示例代码:
// Saga处理初始创建订单的消息 public async Task Handle(OrderCreatedMessage message, IMessageHandlerContext context) { // 首次检查待缴费用 await CheckPendingPayments(message, context); } // 处理定时检查的消息 public async Task Handle(CheckPendingPaymentsMessage message, IMessageHandlerContext context) { // 实际检查待缴费用的逻辑,更新Data.IsAllPaymentsDone Data.IsAllPaymentsDone = await CheckIfAllPaymentsCompleted(message.OrderId); if (!Data.IsAllPaymentsDone) { var paymentRequest = new StartPaymentCommand() { OrderId = message.OrderId, CompanyId = message.CompanyId, ItemNumber = message.ItemNumber, DueAmount = message.DueAmount }; // 发送消息给支付处理器 await context.Send(paymentRequest); // 调度24小时后的下一次检查 var checkOptions = new SendOptions(); checkOptions.DelayDeliveryWith(TimeSpan.FromHours(24)); await context.Send(new CheckPendingPaymentsMessage { OrderId = message.OrderId, CompanyId = message.CompanyId, ItemNumber = message.ItemNumber, DueAmount = message.DueAmount }, checkOptions); } else { // 所有费用已缴清,结束Saga MarkAsComplete(); } } // 辅助方法:检查是否所有费用已缴清 private async Task<bool> CheckIfAllPaymentsCompleted(Guid orderId) { // 调用业务逻辑查询待缴费用 // 示例:return await paymentRepository.HasPendingPayments(orderId); }
3. 关键优化点说明
- 避免阻塞线程:利用NServiceBus的延迟消息机制,让Saga在后台等待,不会占用线程资源;
- 符合业务需求:实现每日检查的逻辑,而非无意义的20分钟循环发送;
- 内存资源控制:不再生成无限多的消息对象,仅在每次检查后决定是否调度下一次任务,内存占用保持稳定;
- 明确Saga生命周期:当所有费用缴清时调用
MarkAsComplete()结束Saga,避免资源泄漏。
内容的提问来源于stack exchange,提问作者Raju Donthula
相关产品推荐
相关产品推荐

