如何让Rebus Saga模式代码等待所有Saga实例完成后返回结果
问题
我正在使用Rebus框架实现Saga模式,对列表中的每笔支付执行分布式事务。由于Rebus的异步执行特性,当前代码会直接返回结果,即便所有Saga实例尚未执行完毕。我尝试用Task.WhenAll解决,但问题仍未得到解决。请问如何修改代码,使其等待所有Saga实例执行完成后再返回结果?
相关初始代码:
public class CreatePrelimnaryMoneyTransfersCommandHandler : IRequestHandler<CreatePrelimnaryMoneyTransfersCommand, Result<bool>> { private readonly IBus _bus; private readonly IPaymentRepository _paymentRepository; public CreatePrelimnaryMoneyTransfersCommandHandler(IBus bus, IPaymentRepository paymentRepository) { _bus = bus; _paymentRepository = paymentRepository; } public async Task<Result<bool>> Handle(CreatePrelimnaryMoneyTransfersCommand request, CancellationToken cancellationToken) { try { var payments = await _paymentRepository.GetPaymentsByQueryParameters(new PaymentQueryParameters { ShowFailedOrChargedBackWithNoPayoutId = true }); foreach (var payment in payments) { await _bus.Send(new ExecuteCreatePrelimnaryMoneyTransfersSaga(p.PaymentId)); } // 这段代码应该等待所有支付处理完成后再执行 return new Result<bool>(true, string.Empty); } catch (Exception ex) { // 根据需要处理或记录异常 return new Result<bool>(false, ex.Message); } } }
使用Task.WhenAll的尝试代码:
public async Task<Result<bool>> Handle(CreatePrelimnaryMoneyTransfersCommand request, CancellationToken cancellationToken) { var payments = await _paymentRepository.GetPaymentsByQueryParameters(new PaymentQueryParameters { ShowPrelimnaryDeservedPayments = true, LoanId = request.LoanId }); await Task.WhenAll(payments.Select(p => _bus.Send(new ExecuteCreatePrelimnaryMoneyTransfersSaga(p.PaymentId)))); _logger.LogInformation($"我希望这条日志在所有Saga完成后才出现"); return new Result<bool>(true, "已保存失败的支付记录"); }
解决方案
核心原因
Task.WhenAll无效的本质是:Rebus的Send方法仅完成消息投递到队列的操作,不会等待对应的Saga实例执行完业务逻辑。要实现等待所有Saga结束的逻辑,需要让每个Saga在完成工作后主动发送完成通知,再在命令处理器中监听并等待所有通知返回。
具体实现步骤
- 定义Saga完成通知消息
public class PrelimnaryMoneyTransferCompleted { public Guid PaymentId { get; set; } public bool Success { get; set; } public string Message { get; set; } }
- 修改Saga类,执行完成后发送回复通知
在ExecuteCreatePrelimnaryMoneyTransfersSaga对应的Saga业务逻辑中,完成事务后发送回复:
public class PrelimnaryMoneyTransferSaga : Saga<PrelimnaryMoneyTransferSagaData>, IAmInitiatedBy<ExecuteCreatePrelimnaryMoneyTransfersSaga> { private readonly IBus _bus; public PrelimnaryMoneyTransferSaga(IBus bus) { _bus = bus; } public async Task Handle(ExecuteCreatePrelimnaryMoneyTransfersSaga message, CancellationToken cancellationToken) { try { // 执行你的分布式事务核心逻辑... // 逻辑完成后发送成功通知 await _bus.Reply(new PrelimnaryMoneyTransferCompleted { PaymentId = message.PaymentId, Success = true, Message = "处理成功" }); } catch (Exception ex) { // 异常时发送失败通知 await _bus.Reply(new PrelimnaryMoneyTransferCompleted { PaymentId = message.PaymentId, Success = false, Message = ex.Message }); throw; } } // Saga数据映射、配置等其他代码... }
- 修改命令处理器,用
Request代替Send并等待所有回复
将原来的Send改为Request(该方法会等待Saga的回复),再通过Task.WhenAll等待所有Saga执行完成:
public class CreatePrelimnaryMoneyTransfersCommandHandler : IRequestHandler<CreatePrelimnaryMoneyTransfersCommand, Result<bool>> { private readonly IBus _bus; private readonly IPaymentRepository _paymentRepository; private readonly ILogger<CreatePrelimnaryMoneyTransfersCommandHandler> _logger; public CreatePrelimnaryMoneyTransfersCommandHandler(IBus bus, IPaymentRepository paymentRepository, ILogger<CreatePrelimnaryMoneyTransfersCommandHandler> logger) { _bus = bus; _paymentRepository = paymentRepository; _logger = logger; } public async Task<Result<bool>> Handle(CreatePrelimnaryMoneyTransfersCommand request, CancellationToken cancellationToken) { try { var payments = await _paymentRepository.GetPaymentsByQueryParameters( new PaymentQueryParameters { ShowFailedOrChargedBackWithNoPayoutId = true }); // 为每笔支付发送请求并等待Saga回复 var completionTasks = payments.Select(async p => await _bus.Request<PrelimnaryMoneyTransferCompleted>( new ExecuteCreatePrelimnaryMoneyTransfersSaga(p.PaymentId), cancellationToken)); // 等待所有Saga执行完成并收集结果 var results = await Task.WhenAll(completionTasks); // 检查是否有处理失败的支付 var failedItems = results.Where(r => !r.Success).ToList(); if (failedItems.Any()) { var errorMsg = string.Join("; ", failedItems.Select(r => $"支付ID {r.PaymentId}: {r.Message}")); _logger.LogError("部分支付处理失败: {ErrorMsg}", errorMsg); return new Result<bool>(false, errorMsg); } _logger.LogInformation("所有Saga实例已执行完成"); return new Result<bool>(true, "所有支付处理完成"); } catch (Exception ex) { _logger.LogError(ex, "处理支付批次时发生异常"); return new Result<bool>(false, ex.Message); } } }
注意事项
- 超时配置:如果Saga执行耗时较长,建议在
Request方法中指定超时时间(例如Request(..., timeout: TimeSpan.FromMinutes(5))),避免无限制等待。 - 消息可靠性:确保消息队列开启持久化配置,防止Saga执行过程中消息丢失导致处理器永久阻塞。
- 并发控制:若支付数量极大,可考虑分批处理,避免大量并发请求占用过多系统资源。
内容的提问来源于stack exchange,提问作者Ahmed Elbatt
相关产品推荐
相关产品推荐

