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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 18:42:03