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
相关产品推荐
相关产品推荐

