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

如何使用MassTransit Saga配合RabbitMQ实现事件顺序自动执行

实现方案:MassTransit Saga 链式顺序执行事件

以下是基于你当前技术栈实现多事件自动按序触发的完整方案:


核心逻辑

通过Saga状态机管控流程的执行阶段,每个业务步骤执行完成后发布对应完成事件,Saga监听该事件后流转状态、自动触发下一个步骤的执行,直到所有步骤全部结束,全程仅需客户端触发第一个启动事件。


具体实现步骤

  • 第一步:定义Saga状态、事件/命令契约
    首先定义Saga实例存储的状态字段,以及流程全链路用到的事件、命令结构:
// Saga状态类,存储流程执行的中间数据和当前状态
public class WorkflowSagaState : SagaStateMachineInstance
{
    public Guid CorrelationId { get; set; }
    public int CurrentStep { get; set; }
    // 存储流程全链路的业务数据,可按实际需求扩展字段
    public string? WorkflowBusinessData { get; set; }
    public State? CurrentState { get; set; }
}

// 客户端触发的初始启动事件,仅需调用这一次
public record StartWorkflow(Guid CorrelationId, string InitBusinessData);
// 各步骤执行完成后的通知事件,由对应步骤的消费者发布
public record EventOneCompleted(Guid CorrelationId, string StepOneResult);
public record EventTwoCompleted(Guid CorrelationId, string StepTwoResult);
// 全流程完成的最终通知事件
public record WorkflowAllFinished(Guid CorrelationId, string FinalResult);

// 各步骤的执行命令,Saga会按顺序自动发布
public record ExecuteEventOne(Guid CorrelationId, string BusinessData);
public record ExecuteEventTwo(Guid CorrelationId, string BusinessData);
  • 第二步:编写Saga状态机,定义状态流转和自动触发逻辑
public class WorkflowStateMachine : MassTransitStateMachine<WorkflowSagaState>
{
    public WorkflowStateMachine()
    {
        // 绑定Saga的状态存储字段
        InstanceState(x => x.CurrentState);

        // 配置所有事件和Saga实例的关联规则,通过CorrelationId绑定同一条流程
        Event(() => StartWorkflow, x => x.CorrelateById(context => context.Message.CorrelationId));
        Event(() => EventOneCompleted, x => x.CorrelateById(context => context.Message.CorrelationId));
        Event(() => EventTwoCompleted, x => x.CorrelateById(context => context.Message.CorrelationId));

        // 初始状态:收到客户端的启动事件后,自动触发第一个事件执行
        Initially(
            When(StartWorkflow)
                .Then(context => 
                {
                    context.Instance.WorkflowBusinessData = context.Data.InitBusinessData;
                    context.Instance.CurrentStep = 1;
                })
                // 发送第一个步骤的执行命令,对应消费者执行后会发布完成事件
                .SendAsync(new Uri("exchange:event-one-execute"), context => 
                    new ExecuteEventOne(context.Instance.CorrelationId, context.Instance.WorkflowBusinessData))
                .TransitionTo(WaitingForEventOneDone)
        );

        // 收到EventOne完成事件后,自动触发EventTwo执行
        During(WaitingForEventOneDone,
            When(EventOneCompleted)
                .Then(context => 
                {
                    // 把上一步的执行结果合并,传给下一步使用
                    context.Instance.WorkflowBusinessData += $";Step1Result:{context.Data.StepOneResult}";
                    context.Instance.CurrentStep = 2;
                })
                .SendAsync(new Uri("exchange:event-two-execute"), context => 
                    new ExecuteEventTwo(context.Instance.CorrelationId, context.Instance.WorkflowBusinessData))
                .TransitionTo(WaitingForEventTwoDone)
        );

        // 收到EventTwo完成事件后,流程结束
        During(WaitingForEventTwoDone,
            When(EventTwoCompleted)
                .Then(context => 
                {
                    context.Instance.WorkflowBusinessData += $";Step2Result:{context.Data.StepTwoResult}";
                    context.Instance.CurrentStep = 3;
                })
                // 可选:发布全流程完成事件供其他服务监听
                .Publish(context => new WorkflowAllFinished(context.Instance.CorrelationId, context.Instance.WorkflowBusinessData))
                .Finalize()
        );

        // 配置流程结束后自动清理Saga实例(按需开启)
        SetCompletedWhenFinalized();
    }

    // 定义Saga的中间状态
    public State WaitingForEventOneDone { get; }
    public State WaitingForEventTwoDone { get; }

    // 注册Saga监听的所有事件
    public Event<StartWorkflow> StartWorkflow { get; }
    public Event<EventOneCompleted> EventOneCompleted { get; }
    public Event<EventTwoCompleted> EventTwoCompleted { get; }
}
  • 第三步:实现各步骤的业务消费者,执行完成后发布完成事件
    以EventOne的消费者为例:
public class EventOneConsumer : IConsumer<ExecuteEventOne>
{
    public async Task Consume(ConsumeContext<ExecuteEventOne> context)
    {
        // 执行EventOne对应的业务逻辑
        var stepResult = await DoEventOneBusinessLogic(context.Message.BusinessData);
        
        // 执行完成后发布完成事件,通知Saga流转到下一步
        await context.Publish(new EventOneCompleted(context.Message.CorrelationId, stepResult));
    }
}
  • 第四步:在服务中注册Saga和消费者
    以.NET的Program.cs配置为例:
builder.Services.AddMassTransit(x =>
{
    // 注册Saga状态机,生产环境建议替换为持久化存储(EF Core/Redis等)
    x.AddSagaStateMachine<WorkflowStateMachine, WorkflowSagaState>()
        .InMemoryRepository();
    
    // 注册所有步骤的消费者
    x.AddConsumer<EventOneConsumer>();
    x.AddConsumer<EventTwoConsumer>();
    
    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host("rabbitmq://你的RabbitMQ地址", h =>
        {
            h.Username("账号");
            h.Password("密码");
        });
        // 自动配置所有端点
        cfg.ConfigureEndpoints(context);
    });
});

注意事项

  • 生产环境禁止使用内存Saga存储,服务重启会丢失所有进行中的流程状态,需替换为持久化存储方案
  • 建议为所有消费者添加幂等校验,通过CorrelationId + 步骤号判断当前步骤是否已经执行过,避免重复投递导致的业务异常
  • 可按需添加步骤超时、失败重试、死信队列规则,处理步骤执行失败的异常场景
  • 如果是动态步骤的场景,可把步骤配置存储在Saga状态中,通过通用完成事件处理逻辑动态判断下一步要触发的命令,不需要逐个定义状态

内容的提问来源于stack exchange,提问作者Laith Sa'd Al-Deen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 13:15:04