Mass Transit状态机复合事件无法触发问题求助
MassTransit状态机复合事件AdminsUpdated无法触发排查
我定义了一个基于MassTransit的SubscriptionStateMachine状态机,目前可以正常处理AdminUpdatedEvent和SuperAdminUpdatedEvent,但依赖这两个事件的复合事件AdminsUpdated始终无法触发。已排查SEQ收发日志、数据库状态机数据(is_new_tenant为true,admins_updated_status为2)及容器日志,仍未解决。
状态机代码如下:
public sealed class SubscriptionStateMachine : MassTransitStateMachine<SubscriptionState> { // States public State UpdatingTenantState { get; set; } = null!; public State UpdatingLocalSubscriptionState { get; set; } = null!; public State SendingCommunicationState { get; set; } = null!; public State FailedState { get; set; } = null! // Events public Event<StripeSubscriptionUpdatedEvent> StripeSubscriptionUpdated { get; set; } = null!; public Event<TenantUpdatedEvent> TenantUpdated { get; set; } = null!; public Event<LocalSubscriptionUpdatedEvent> LocalSubscriptionUpdated { get; set; } = null!; public Event<AdminUpdatedEvent> AdminUpdated { get; set; } = null!; public Event<SuperAdminUpdatedEvent> SuperAdminUpdated { get; set; } = null!; public Event<WelcomeEmailSentEvent> WelcomeEmailSent { get; set; } = null!; public Event<SubscriptionUpdatedEmailSentEvent> SubscriptionUpdatedEmailSent { get; set; } = null!; // Faults public Event<Fault<UpdateTenantCommand>> UpdateTenantFaulted { get; set; } = null!; public Event<Fault<UpdateLocalSubscriptionCommand>> UpdateLocalSubscriptionFaulted { get; set; } = null!; public Event<Fault<UpdateAdminCommand>> UpdateAdminFaulted { get; set; } = null!; public Event<Fault<UpdateSuperAdminCommand>> UpdateSuperAdminFaulted { get; set; } = null!; public SubscriptionStateMachine() { InstanceState(x => x.CurrentState); Event(() => StripeSubscriptionUpdated, x => x.CorrelateById(m => m.Message.CorrelationId)); Event(() => TenantUpdated, x => x.CorrelateById(m => m.Message.CorrelationId)); Event(() => LocalSubscriptionUpdated, x => x.CorrelateById(m => m.Message.CorrelationId)); Event(() => AdminUpdated, x => x.CorrelateById(m => m.Message.CorrelationId)); Event(() => SuperAdminUpdated, x => x.CorrelateById(m => m.Message.CorrelationId)); Event(() => WelcomeEmailSent, x => x.CorrelateById(m => m.Message.CorrelationId)); Event(() => SubscriptionUpdatedEmailSent, x => x.CorrelateById(m => m.Message.CorrelationId)); Event(() => UpdateTenantFaulted, x => x.CorrelateById(m => m.Message.Message.CorrelationId)); Event(() => UpdateLocalSubscriptionFaulted, x => x.CorrelateById(m => m.Message.Message.CorrelationId)); Event(() => UpdateAdminFaulted, x => x.CorrelateById(m => m.Message.Message.CorrelationId)); Event(() => UpdateSuperAdminFaulted, x => x.CorrelateById(m => m.Message.Message.CorrelationId)); // Create or update tenant Initially( When(StripeSubscriptionUpdated) .Then(context => { context.Saga.CorrelationId = context.Message.CorrelationId; context.Saga.Identifier = context.Message.Identifier; context.Saga.SubscriptionId = context.Message.SubscriptionId; context.Saga.CustomerEmail = context.Message.CustomerEmail; context.Saga.CustomerName = context.Message.CustomerName; context.Saga.Description = context.Message.Description; context.Saga.CurrentPeriodStart = context.Message.CurrentPeriodStart; context.Saga.CurrentPeriodEnd = context.Message.CurrentPeriodEnd; context.Saga.Created = context.Message.Created; context.Saga.Status = context.Message.Status; }) .SendAsync(new Uri("queue:update-tenant"), context => context.Init<UpdateTenantCommand>(new { CorrelationId = context.Saga.CorrelationId, Identifier = context.Saga.Identifier, Status = context.Saga.Status })) .TransitionTo(UpdatingTenantState)); // Create or update the local subscription During(UpdatingTenantState, When(TenantUpdated) .Then(context => { context.Saga.TenantId = context.Message.TenantId; context.Saga.IsNewTenant = context.Message.IsNewTenant; }).SendAsync(new Uri("queue:update-local-subscription"), context => context.Init<UpdateLocalSubscriptionCommand>(new { CorrelationId = context.Saga.CorrelationId, TenantId = context.Saga.TenantId, Identifier = context.Saga.Identifier, StripeSubscriptionId = context.Saga.SubscriptionId, CustomerName = context.Saga.CustomerName, CustomerEmail = context.Saga.CustomerEmail, Description = context.Saga.Description, Status = context.Saga.Status, CurrentPeriodStart = context.Saga.CurrentPeriodStart, CurrentPeriodEnd = context.Saga.CurrentPeriodEnd, Created = context.Saga.Created })).TransitionTo(UpdatingLocalSubscriptionState)); // Create or update the admins During(UpdatingLocalSubscriptionState, When(LocalSubscriptionUpdated) .SendAsync(new Uri("queue:update-admin"), p => p.Init<UpdateAdminCommand>(new { CorrelationId = p.Saga.CorrelationId, TenantId = p.Saga.TenantId, Identifier = p.Saga.Identifier, CustomerEmail = p.Saga.CustomerEmail, })) .SendAsync(new Uri("queue:update-super-admin"), p => p.Init<UpdateSuperAdminCommand>(new { CorrelationId = p.Saga.CorrelationId, TenantId = p.Saga.TenantId, Identifier = p.Saga.Identifier, }))); // Composite event once the tenant admins are updated CompositeEvent( () => AdminsUpdated, x => x.AdminsUpdatedStatus, AdminUpdated, SuperAdminUpdated); DuringAny( When(AdminsUpdated) .If(context => context.Saga.IsNewTenant, then => then.SendAsync(new Uri("queue:send-welcome-email"), s => s.Init<SendWelcomeEmailCommand>(new { CorrelationId = s.Saga.CorrelationId, TenantId = s.Saga.TenantId, Identifier = s.Saga.Identifier, CustomerEmail = s.Saga.CustomerEmail, CustomerName = s.Saga.CustomerName, CurrentPeriodStart = s.Saga.CurrentPeriodStart, CurrentPeriodEnd = s.Saga.CurrentPeriodEnd, })).TransitionTo(SendingCommunicationState))); DuringAny( When(AdminsUpdated) .If(context => !context.Saga.IsNewTenant, then => then.SendAsync(new Uri("queue:send-subscription-updated-email"), s => s.Init<SendSubscriptionUpdatedEmailCommand>(new { CorrelationId = s.Saga.CorrelationId, Identifier = s.Saga.Identifier, CustomerEmail = s.Saga.CustomerEmail, CustomerName = s.Saga.CustomerName, CurrentPeriodStart = s.Saga.CurrentPeriodStart, CurrentPeriodEnd = s.Saga.CurrentPeriodEnd, })).TransitionTo(SendingCommunicationState))); During(SendingCommunicationState, When(WelcomeEmailSent) .Finalize()); During(SendingCommunicationState, When(SubscriptionUpdatedEmailSent) .Finalize()); // Faults DuringAny(When(UpdateAdminFaulted) .TransitionTo(FailedState)); DuringAny(When(UpdateSuperAdminFaulted) .TransitionTo(FailedState)); DuringAny(When(UpdateTenantFaulted) .TransitionTo(FailedState)); DuringAny(When(UpdateLocalSubscriptionFaulted) .TransitionTo(FailedState)); SetCompletedWhenFinalized(); } // Composite events public Event AdminsUpdated { get; set; } = null!; }
排查方向建议
- 复合事件配置顺序:将
CompositeEvent的注册代码移到所有状态转换逻辑之前(比如Initially代码块之前),MassTransit需要先识别复合事件定义,才能正确跟踪依赖事件的触发状态。 - 事件关联验证:确认
AdminUpdated和SuperAdminUpdated事件的CorrelationId与状态机实例的CorrelationId完全匹配,包括字符串格式(如GUID的连字符、大小写)。 - 状态跟踪字段类型:检查
SubscriptionState中的AdminsUpdatedStatus字段是否为int类型,MassTransit用位掩码(整数)跟踪复合事件的触发状态,非整数类型会导致状态无法正确累加。 - 状态机状态检查:当
AdminUpdated和SuperAdminUpdated触发时,查看SEQ日志中状态机实例的CurrentState值,确认实例是否处于允许接收这两个事件的状态。 - 复合事件触发逻辑:默认情况下复合事件仅触发一次,如果之前已触发但未正确处理,检查状态机实例是否被
SetCompletedWhenFinalized()标记为完成,导致无法再响应事件。
内容的提问来源于stack exchange,提问作者user351479
相关产品推荐
相关产品推荐

