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

如何在MassTransit Saga状态机的Activity中触发定义的事件?

问题:在MassTransit Saga Activity中触发状态机定义的事件

我正在构建一个Saga状态机,当前实现中,状态机定义了CaseCreationFinished和CaseCreationFailed两个事件,但在Activity内部无法访问这些事件来通过context.Raise触发。以下是我的精简代码:

状态机代码

public class DueDiligenceCaseCreateStateMachine : MassTransitStateMachine<DueDiligenceCaseCreateState>
{
    public State CreatingCase { get; private set; }

    public Event<DueDiligenceCaseCreateCommand> TriggerReceived { get; private set; }

    public Event CaseCreationFinished { get; private set; }

    public Event CaseCreationFailed { get; private set; }

    private readonly ILogger<DueDiligenceCaseCreateStateMachine> _logger;
    private readonly IOptions<DueDiligenceCaseCreateSagaOptions> _sagaOptions;

    public DueDiligenceCaseCreateStateMachine(
        ILogger<DueDiligenceCaseCreateStateMachine> logger,
        IOptions<DueDiligenceCaseCreateSagaOptions> sagaOptions)
    {
        _logger = logger;
        _sagaOptions = sagaOptions;

        Configure();
        BuildProcess();
    }

    private void Configure()
    {
        Event(
            () => TriggerReceived,
            e => e.CorrelateById(x => x.Message.DueDiligenceCaseId));
    }

    private void BuildProcess()
    {
        During(
            Initial,
            When(TriggerReceived)
                .TransitionTo(CreatingCase)
                .Activity(CreateCase));
    }
    
    
    private EventActivityBinder<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> CreateCase(
        IStateMachineActivitySelector<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> sel) =>
        sel.OfType<CreateCaseActivity>();
}

Activity代码

public class CreateCaseActivity : BaseActivity<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand>
{
    private readonly ICommandHandler<InitializeCaseCommand> _initializeCaseHandler;
    private readonly IOptions<ApplicationOptions> _options;
    private readonly ILogger<DueDiligenceCaseCreateConsumer> _logger;

    public CreateCaseActivity(
        ICommandHandler<InitializeCaseCommand> initializeCaseHandler,
        IOptions<ApplicationOptions> options,
        ILogger<DueDiligenceCaseCreateConsumer> logger)
    {
        _initializeCaseHandler = initializeCaseHandler;
        _options = options;
        _logger = logger;
    }

    public override async Task Execute(
        BehaviorContext<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> context,
        Behavior<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> next)
    {
        _logger.LogInformation(
            "Consuming {Command} started, case id: {caseid}, creating a case...",
            nameof(DueDiligenceCaseCreateCommand),
            context.Data.DueDiligenceCaseId);

        var initializeCaseCmd = ConvertMessageToCommand(context.Data);
        
        initializeCaseCmd.CanHaveOnlyOneActiveCasePerCustomer = !_options.Value.FeatureToggles.AllowMultipleActiveCasesOnSingleCustomer;

        try
        {
            await _initializeCaseHandler.Handle(initializeCaseCmd);
        }
        catch 
        {
        
        }
        finally
        {
            await next.Execute(context);
        }
    }

    private InitializeCaseCommand ConvertMessageToCommand(DueDiligenceCaseCreateCommand message) =>
        // returns the command object
}

可行的解决办法

方法1:直接注入状态机事件到Activity

状态机的事件是公开属性,可以通过构造函数直接注入到Activity中,这样就能在Activity内部调用context.Raise触发事件:

修改后的Activity代码

public class CreateCaseActivity : BaseActivity<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand>
{
    private readonly ICommandHandler<InitializeCaseCommand> _initializeCaseHandler;
    private readonly IOptions<ApplicationOptions> _options;
    private readonly ILogger<DueDiligenceCaseCreateConsumer> _logger;
    private readonly Event _caseCreationFinished;
    private readonly Event _caseCreationFailed;

    public CreateCaseActivity(
        ICommandHandler<InitializeCaseCommand> initializeCaseHandler,
        IOptions<ApplicationOptions> options,
        ILogger<DueDiligenceCaseCreateConsumer> logger,
        Event caseCreationFinished,
        Event caseCreationFailed)
    {
        _initializeCaseHandler = initializeCaseHandler;
        _options = options;
        _logger = logger;
        _caseCreationFinished = caseCreationFinished;
        _caseCreationFailed = caseCreationFailed;
    }

    public override async Task Execute(
        BehaviorContext<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> context,
        Behavior<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> next)
    {
        _logger.LogInformation(
            "Consuming {Command} started, case id: {caseid}, creating a case...",
            nameof(DueDiligenceCaseCreateCommand),
            context.Data.DueDiligenceCaseId);

        var initializeCaseCmd = ConvertMessageToCommand(context.Data);
        
        initializeCaseCmd.CanHaveOnlyOneActiveCasePerCustomer = !_options.Value.FeatureToggles.AllowMultipleActiveCasesOnSingleCustomer;

        try
        {
            await _initializeCaseHandler.Handle(initializeCaseCmd);
            // 触发成功事件
            await context.Raise(_caseCreationFinished);
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "Case creation failed for id: {caseid}", context.Data.DueDiligenceCaseId);
            // 触发失败事件
            await context.Raise(_caseCreationFailed);
        }
        finally
        {
            await next.Execute(context);
        }
    }

    private InitializeCaseCommand ConvertMessageToCommand(DueDiligenceCaseCreateCommand message) =>
        new InitializeCaseCommand { /* 这里补充你的转换逻辑 */ };
}

在状态机中传递事件

修改状态机中注册Activity的代码,把状态机的事件传递进去:

private EventActivityBinder<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> CreateCase(
    IStateMachineActivitySelector<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> sel) =>
    sel.OfType<CreateCaseActivity>(() => new CreateCaseActivity(
        _serviceProvider.GetRequiredService<ICommandHandler<InitializeCaseCommand>>(),
        _serviceProvider.GetRequiredService<IOptions<ApplicationOptions>>(),
        _serviceProvider.GetRequiredService<ILogger<DueDiligenceCaseCreateConsumer>>(),
        CaseCreationFinished,
        CaseCreationFailed));

注:需要确保状态机能拿到IServiceProvider实例,或者直接通过DI容器注册Activity时注入事件。

方法2:通过状态机访问器获取事件

注入IStateMachineAccessor<DueDiligenceCaseCreateState>到Activity,通过它拿到状态机实例,进而访问事件:

public class CreateCaseActivity : BaseActivity<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand>
{
    private readonly ICommandHandler<InitializeCaseCommand> _initializeCaseHandler;
    private readonly IOptions<ApplicationOptions> _options;
    private readonly ILogger<DueDiligenceCaseCreateConsumer> _logger;
    private readonly DueDiligenceCaseCreateStateMachine _stateMachine;

    public CreateCaseActivity(
        ICommandHandler<InitializeCaseCommand> initializeCaseHandler,
        IOptions<ApplicationOptions> options,
        ILogger<DueDiligenceCaseCreateConsumer> logger,
        IStateMachineAccessor<DueDiligenceCaseCreateState> stateMachineAccessor)
    {
        _initializeCaseHandler = initializeCaseHandler;
        _options = options;
        _logger = logger;
        _stateMachine = (DueDiligenceCaseCreateStateMachine)stateMachineAccessor.StateMachine;
    }

    public override async Task Execute(
        BehaviorContext<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> context,
        Behavior<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> next)
    {
        // 原有业务逻辑...
        try
        {
            await _initializeCaseHandler.Handle(initializeCaseCmd);
            await context.Raise(_stateMachine.CaseCreationFinished);
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "Case creation failed for id: {caseid}", context.Data.DueDiligenceCaseId);
            await context.Raise(_stateMachine.CaseCreationFailed);
        }
        finally
        {
            await next.Execute(context);
        }
    }
}

方法3:状态机集中处理事件触发(推荐)

遵循Saga关注点分离原则,Activity只负责执行业务逻辑并将结果存入Saga状态,由状态机统一处理事件触发:

1. 修改Saga状态类,添加结果标记

public class DueDiligenceCaseCreateState : SagaStateMachineInstance
{
    public Guid CorrelationId { get; set; }
    public string CurrentState { get; set; }
    // 添加标记字段
    public bool CaseCreationSucceeded { get; set; }
}

2. 修改Activity,存储执行结果

public override async Task Execute(
    BehaviorContext<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> context,
    Behavior<DueDiligenceCaseCreateState, DueDiligenceCaseCreateCommand> next)
{
    bool creationSucceeded = false;
    try
    {
        await _initializeCaseHandler.Handle(initializeCaseCmd);
        creationSucceeded = true;
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "Case creation failed for id: {caseid}", context.Data.DueDiligenceCaseId);
        creationSucceeded = false;
    }
    finally
    {
        // 将结果存入Saga状态
        context.Instance.CaseCreationSucceeded = creationSucceeded;
        await next.Execute(context);
    }
}

3. 状态机中添加事件触发逻辑

private void BuildProcess()
{
    During(Initial,
        When(TriggerReceived)
            .TransitionTo(CreatingCase)
            .Activity(CreateCase)
            // 根据状态结果触发对应事件
            .Then(context => 
            {
                if (context.Instance.CaseCreationSucceeded)
                    context.Raise(CaseCreationFinished);
                else
                    context.Raise(CaseCreationFailed);
            }));
}

内容的提问来源于stack exchange,提问作者Marek M.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 00:10:47