MassTransit状态机处理响应后仍收到RequestTimeoutExpired事件问题
我实现了一个MassTransit状态机,流程为:接收初始事件后发送请求,等待响应并转换至终态。正常情况下功能符合预期,但出现异常情况:状态机已成功处理响应、更新Saga并切换至Final状态后,仍会收到RequestTimeoutExpired事件,抛出如下异常:
MassTransit.NotAcceptedStateMachineException: MassTransitExample.Models.Book(a16bd901-3ed9-4633-8e0c-4b00d8babeff) Saga exception on receipt of MassTransit.Contracts.RequestTimeoutExpired<MassTransitExample.Messages.TestRequest>: Not accepted in state Final
状态机实现
public class BookStateMachine : MassTransitStateMachine<Book> { public BookStateMachine() { Event(() => Added, x => x.CorrelateById(m => m.Message.BookId)); Request(() => TestRequest, r => { r.Completed = e => e.ConfigureConsumeTopology = false; r.Faulted = e => e.ConfigureConsumeTopology = false; //r.TimeoutExpired = e => e.ConfigureConsumeTopology = false; }); InstanceState(x => x.CurrentState, Available); Initially( When(Added) .CopyDataToInstance() .Request(TestRequest, ctx => new TestRequest() { BookId = ctx.Message.BookId }) .TransitionTo(TestRequest.Pending)); this.During(TestRequest.Pending, When(TestRequest.Completed) .Then(ctx => { ctx.Saga.ResponseText = ctx.Message.ResponseText; }) .TransitionTo(Final)); } public Event<BookAdded> Added { get; private set; } public Request<Book, TestRequest, TestResponse> TestRequest { get; } = null!; public State Available { get; private set; } } public static class BookStateMachineExtensions { public static EventActivityBinder<Book, BookAdded> CopyDataToInstance( this EventActivityBinder<Book, BookAdded> binder) { return binder .Then(x => { x.Saga.BookId = x.Message.BookId; x.Saga.DateAdded = x.Message.Timestamp.Date; x.Saga.Title = x.Message.Title; x.Saga.Isbn = x.Message.Isbn; }); } }
DbContext实现
public class TestDbContext(DbContextOptions options) : SagaDbContext(options) { public DbSet<Book> Books { get; set; } protected override IEnumerable<ISagaClassMap> Configurations { get { yield return new BookStateMap(); } } protected override void OnModelCreating(ModelBuilder modelBuilder) { foreach (ISagaClassMap sagaMap in Configurations) { sagaMap.Configure(modelBuilder); } } } public class BookStateMap : SagaClassMap<Book> { protected override void Configure(EntityTypeBuilder<Book> entity, ModelBuilder model) { entity.HasKey(x => x.CorrelationId); entity.Property(x => x.CurrentState).HasMaxLength(64); entity.Property(x => x.BookId).IsRequired(); entity.Property(x => x.DateAdded).IsRequired(); entity.Property(x => x.Title).IsRequired(); entity.Property(x => x.Isbn).IsRequired(); entity.Property(x => x.ResponseText); entity.Property(x => x.CurrentState).IsRequired(); entity.Property(x => x.RowVersion).IsRowVersion(); } }
服务配置
builder.Services.AddDbContext<TestDbContext>( options => options.UseSqlServer(connectionString)); builder.Services.AddMassTransit(x => { x.AddDelayedMessageScheduler(); Assembly? entryAssembly = Assembly.GetEntryAssembly(); x.AddConsumers(entryAssembly); x.AddSagaStateMachine<BookStateMachine, Book>() .EntityFrameworkRepository(r => { r.ConcurrencyMode = ConcurrencyMode.Optimistic; r.ExistingDbContext<TestDbContext>(); r.UseSqlServer(); }); x.UsingRabbitMq((context, cfg) => { cfg.UseDelayedMessageScheduler(); cfg.Host("localhost", "/", h => { h.Username("guest"); h.Password("guest"); }); cfg.ConfigureEndpoints(context); }); });
消费者与消息契约
public class TestRequestConsumer : IConsumer<TestRequest> { public async Task Consume(ConsumeContext<TestRequest> context) { await Task.Delay(3000); await context.RespondAsync(new TestResponse() { BookId = context.Message.BookId, ResponseText = "Responded" }); } } public record BookAdded { public Guid BookId { get; init; } public string Title { get; init; } public string Isbn { get; init; } public DateTime Timestamp { get; init; } } public record TestRequest { public Guid BookId { get; init; } } public record TestResponse { public Guid BookId { get; init; } public string ResponseText { get; init; } }
解决方案
这个问题的核心是:当Saga进入Final状态后,之前发起的Request的超时事件仍会被路由到Saga,但Final状态默认不处理任何未明确配置的事件,因此抛出异常。
方法一:阻止超时事件路由到Saga
在Request配置中取消注释TimeoutExpired的拓扑配置,让MassTransit不为该事件创建消费拓扑,这样超时事件就不会被发送到Saga的队列:
Request(() => TestRequest, r => { r.Completed = e => e.ConfigureConsumeTopology = false; r.Faulted = e => e.ConfigureConsumeTopology = false; r.TimeoutExpired = e => e.ConfigureConsumeTopology = false; // 启用这行 });
方法二:让Final状态忽略超时事件
如果需要保留超时事件的拓扑(比如后续有其他处理逻辑),可以在状态机中配置Final状态明确忽略TestRequest.TimeoutExpired事件:
public BookStateMachine() { // ... 其他现有代码 ... During(Final, When(TestRequest.TimeoutExpired) .Ignore()); }
额外说明
当Saga进入Final状态后,MassTransit默认会保留Saga实例(除非配置了自动删除),但此时它不再处理未明确绑定的事件。通过上述两种方式,要么从根源阻止超时事件到达Saga,要么让Final状态明确忽略该事件,就能解决异常问题。
内容的提问来源于stack exchange,提问作者Sugafree

