MassTransit中CompositeEvent未触发且ProcessingStatus值异常问题
问题根因
复合事件不触发、ProcessingStatus值异常是三个问题叠加导致的:
- 复合事件配置顺序错误:MassTransit状态机的行为注册是顺序生效的,你在已经注册完
Initial、Processing状态下LinkContract、LinkCustomer两个事件的处理逻辑之后,才声明CompositeEvent,导致复合事件的位标记更新逻辑根本没有被插入到两个事件的处理管道中,事件到达时不会自动更新ProcessingStatus的位值。 - 误用
CompositeEventOptions.IncludeInitial配置:该选项的作用是将「Saga处于初始未流转状态」作为复合事件触发的第三个必要条件,你当前逻辑收到第一个事件就会流转到Processing状态,永远无法满足三个条件同时成立的触发要求。 - Redis持久化序列化失败:你定义的
ProcessingStatus是CompositeEventStatus值类型,默认的JSON序列化器无法正确识别该结构体的内部_value字段,第一个事件处理完写入Redis时,ProcessingStatus的值没有被正确持久化,第二个事件读取Saga实例时该字段被反序列化为默认值0,完全丢失了之前的位标记。
修复方案
- 调整状态机配置顺序,先声明所有事件、注册复合事件,再编写各状态下的事件处理逻辑,保证复合事件的位更新逻辑能被正确注入到事件处理管道最前端。
- 移除不需要的
CompositeEventOptions.IncludeInitial参数,该场景下不需要将初始状态作为触发条件。 - 将
BillingState中ProcessingStatus的类型从CompositeEventStatus改为int,CompositeEventStatus与int支持隐式互转,MassTransit完全兼容int类型的复合事件状态跟踪字段,可以彻底规避结构体序列化失败的问题。
修正后的代码
首先修正状态类定义:
using MassTransit; namespace Sandbox.WebApi.StateMachines; public class BillingState : SagaStateMachineInstance, ISagaVersion { public Guid CorrelationId { get; set; } public string CurrentState { get; set; } // 改为int类型,规避结构体序列化问题 public int ProcessingStatus { get; set; } public int Version { get; set; } public Guid CustomerId { get; set; } public Guid ContractId { get; set; } }
然后修正状态机配置:
using MassTransit; using Sandbox.WebApi.Models.StateEvents; using Serilog; namespace Sandbox.WebApi.StateMachines; public class BillingStateMachine : MassTransitStateMachine<BillingState> { public BillingStateMachine() { InstanceState(x => x.CurrentState); // 第一步:声明所有事件 Event(() => LinkCustomer, e => e .CorrelateById(x => x.Message.CustomerId)); Event(() => LinkContract, e => e .CorrelateById(x => x.Message.CustomerId)); // 第二步:提前注册复合事件,移除错误的IncludeInitial配置 CompositeEvent(() => PaymentSucceeded, x => x.ProcessingStatus, CompositeEventOptions.None, LinkContract, LinkCustomer); // 第三步:注册各状态下的事件处理逻辑 Initially( When(LinkCustomer) .Then(HandleLinkCustomer) .TransitionTo(Processing), When(LinkContract) .Then(HandleLinkContract) .TransitionTo(Processing)); During(Processing, When(LinkContract) .Then(HandleLinkContract), When(LinkCustomer) .Then(HandleLinkCustomer), When(PaymentSucceeded) .Then(HandlePaymentSucceeded) .TransitionTo(Paid)); } public State Processing { get; } public State Paid { get; } public Event PaymentSucceeded { get; } public Event<ContractCreated> LinkContract { get; } public Event<CustomerCreated> LinkCustomer { get; } private static void HandleLinkCustomer(BehaviorContext<BillingState, CustomerCreated> context) { context.Saga.CustomerId = context.Message.CustomerId; Log.Information("Customer was linked: {CustomerId}", context.Saga.CustomerId); } private static void HandleLinkContract(BehaviorContext<BillingState, ContractCreated> context) { context.Saga.ContractId = context.Message.ContractId; context.Saga.CustomerId = context.Message.CustomerId; Log.Information("Contract was linked: {ContractId}", context.Saga.ContractId); } private static void HandlePaymentSucceeded(BehaviorContext<BillingState> context) { Log.Information("Billing was paid: {ContractId} & {CustomerId}", context.Saga.ContractId, context.Saga.CustomerId); } }
依赖注入配置不需要修改即可正常运行。
修复后预期行为
- 收到第一个事件(
LinkContract/LinkCustomer任意一个)时,ProcessingStatus被置为1或2,Saga流转到Processing状态,值被正确持久化到Redis。 - 收到第二个事件时,位运算后
ProcessingStatus值为3,复合事件PaymentSucceeded自动触发,执行支付成功逻辑,Saga流转到Paid状态。
内容的提问来源于stack exchange,提问作者pilkha
相关产品推荐
相关产品推荐

