如何在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.
相关产品推荐
相关产品推荐

