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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 17:38:10