如何使用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
相关产品推荐
相关产品推荐

