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

MassTransit状态机Saga无法通过Send消息触发,仅Publish可用的问题排查求助

MassTransit状态机Saga无法通过Send消息触发,仅Publish可用的问题排查求助

我最近在配置MassTransit的状态机Saga时遇到了一个棘手的问题:当我通过Send方法发送InitiateAccountDeletion消息时,消息一直卡在对应的队列里,完全没法触发Saga启动;但如果换成Publish方法发送,Saga就能正常工作。我认为这个业务场景下用Send才是正确的做法,所以想请教大家我到底遗漏了什么配置?

我强烈怀疑InitiateAccountDeletion队列的消费者端点没有被正确配置,但我本来以为ConfigureEndpoints()方法会自动完成这件事。

以下是我的相关代码:

Saga实例定义

public class DeleteAccountSaga : SagaStateMachineInstance
{
    public Guid CorrelationId { get; set; }
    public string? CurrentState { get; set; }
    public int FilesState { get; set; }
}

状态机逻辑

public class DeleteAccountStateMachine : MassTransitStateMachine<DeleteAccountSaga>
{
    // Files are being deleted.
    public State DeletingFiles { get; }
    // Files were deleted.
    public State FilesDeletionCompleted { get; }
    // The account and related data is being deleted from the database.
    public State DeletingAccount { get; }
    // The account and related data have been deleted from the database.
    public State AccountDeletionCompleted { get; }

    // This message should initiate the saga.
    public Event<InitiateAccountDeletion> InitiateAccountDeletion { get; }
    // This event indicates that files from a specific server (T1) were deleted.
    public Event<AccountFilesT1Deleted> AccountFilesT1Deleted { get; }
    // This event indicates that files from a specific server (T2) were deleted.
    public Event<AccountFilesT2Deleted> AccountFilesT2Deleted { get; }
    // Composite event that indicates that files from all servers were deleted.
    public Event AccountFilesDeleted { get; }
    // This event indicates that the database records were deleted.
    public Event<AccountDeleted> AccountDeleted { get; }

    public DeleteAccountStateMachine()
    {
        InstanceState(x => x.CurrentState);

        Event(() => InitiateAccountDeletion, x =>
        {
            x.InsertOnInitial = true;
            x.CorrelateById(context => context.Message.AccountId);
            x.SetSagaFactory(context => new()
            {
                CorrelationId = context.Message.AccountId
            });
        });

        Event(() => AccountFilesT1Deleted, x => x.CorrelateById(context => context.Message.AccountId));
        Event(() => AccountFilesT2Deleted, x => x.CorrelateById(context => context.Message.AccountId));
        Event(() => AccountDeleted, x => x.CorrelateById(context => context.Message.AccountId));

        // The saga should be started when the InitiateAccountDeletion message is sent.
        // After, it should send DeleteFilesT1 and DeleteFilesT2 messages so that files from corresponsing servers are deleted.
        Initially(
            When(InitiateAccountDeletion)
                .Send(context => new DeleteFilesT1(context.Message.AccountId))
                .Send(context => new DeleteFilesT2(context.Message.AccountId))
                .TransitionTo(DeletingFiles)
        );

        // Indicate that all files were deleted.
        CompositeEvent(() => FilesDeleted, x => x.FilesState,
            AccountFilesT1Deleted,
            AccountFilesT2Deleted
        );

        // When all files have been deleted, send DeleteAccount message which will delete all data from the database.
        During(DeletingFiles,
            Ignore(InitiateAccountDeletion),
            When(FilesDeleted)
                .TransitionTo(FilesDeletionCompleted)
                .Send(context => new DeleteAccount(context.Saga.CorrelationId))
        );

        During(FilesDeletionCompleted,
            When(AccountDeleted)
                .TransitionTo(AccountDeletionCompleted)
                .Finalize()
        );

        SetCompletedWhenFinalized();
    }
}

Startup配置

public class Startup
{
    public void ConfigureServices(IServiceCollection services)
    {
        // [...]
        EndpointConvention.Map<InitiateAccountDeletion>(new Uri($"queue:{nameof(InitiateAccountDeletion)}"));
        EndpointConvention.Map<DeleteFilesT1>(new Uri($"queue:{nameof(DeleteFilesT1)}"));
        EndpointConvention.Map<DeleteFilesT2>(new Uri($"queue:{nameof(DeleteFilesT2)}"));
        EndpointConvention.Map<DeleteAccount>(new Uri($"queue:{nameof(DeleteAccount)}"));

        services.AddMassTransit(x =>
        {
            x.AddEntityFrameworkOutbox<AccountDbContext>(options =>
            {
                options.UsePostgres();
                options.UseBusOutbox();
            });

            // This registers:
            // - DeleteFilesT1Consumer
            // - DeleteFilesT2Consumer
            // - DeleteAccountConsumer
            x.AddConsumers(typeof(Startup).Assembly);

            x.AddSagaStateMachine<DeleteAccountStateMachine, DeleteAccountSaga>()
                .EntityFrameworkRepository(options =>
                {
                    options.ExistingDbContext<AccountDbContext>();
                    options.UsePostgres();
                });

            x.UsingAmazonSqs((context, cfg) =>
            {
                cfg.Host(configuration["AWS:Region"], (_) => { });
                cfg.ConfigureEndpoints(context);
            });
        });
    }
}

控制器发送消息代码

public class AccountController : ControllerBase
{
    [HttpDelete("{accountId}")]
    public async Task<IActionResult> Delete(Guid accountId)
    {
        var account = await _dbContext.Account.FindAsync(accountId);
        if (account == null)
            return NotFound("...");

        account.Status = "Deleting";

        // This message remains in the InitiateAccountDeletion queue. It doesn't get processed by the StateMachine.
        await _sendEndpointProvider.Send<InitiateAccountDeletion>(new(
            AccountId: account.AccountId
        ));

        await _dbContext.SaveChangesAsync();
        return Accepted();
    }
}

Update 1
我修改了一行端点约定的代码后,问题就解决了:
原来的代码:

EndpointConvention.Map<InitiateAccountDeletion>(new Uri($"queue:{nameof(InitiateAccountDeletion)}"));

改成:

EndpointConvention.Map<InitiateAccountDeletion>(new Uri($"queue:{nameof(DeleteTenant)}"));

备注:内容来源于stack exchange,提问作者GomalId

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.22 11:52:58