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

如何让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在完成工作后主动发送完成通知,再在命令处理器中监听并等待所有通知返回。

具体实现步骤

  1. 定义Saga完成通知消息
public class PrelimnaryMoneyTransferCompleted
{
    public Guid PaymentId { get; set; }
    public bool Success { get; set; }
    public string Message { get; set; }
}
  1. 修改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数据映射、配置等其他代码...
}
  1. 修改命令处理器,用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:02:09