MassTransit Saga状态机接收Kafka响应后偶发事件不触发
MassTransit Saga状态机结合Kafka偶发事件未触发问题排查求助
使用MassTransit Saga状态机与Kafka编排数据流时,遇到偶发异常:接收Kafka主题的ProducerValidationResult响应后,状态机对应的事件未触发。
Saga配置
private static void RiderConfiguration(IRiderRegistrationConfigurator config, Configuration configuration) { config.AddInMemoryInboxOutbox(); IMongoDatabase db = AddMongoSagaDatabase(configuration.MongoConfiguration); config.AddSagaStateMachine<PipelineStateMachine, PipelineRunState>() .MongoDbRepository(db, cfg => { cfg.DatabaseName = configuration.MongoConfiguration.DatabaseName; }); }
端点配置
private static IKafkaFactoryConfigurator AddConsumerToBuilder<TMessage>( IKafkaFactoryConfigurator configurator, IRiderRegistrationContext context, string topicName, string groupId, Configuration configuration) where TMessage : class { configurator.TopicEndpoint<string, TMessage>( topicName, new ConsumerConfig() { GroupId = groupId, AutoOffsetReset = AutoOffsetReset.Latest }, endpointConfig => { endpointConfig.SetValueDeserializer(new JsonDeserializer<TMessage>()); endpointConfig.UseMessageRetry(r => { r.Immediate(configuration.KafkaConfiguration.RetryCount); r.Handle<MongoDbConcurrencyException>(); }); endpointConfig.UseInMemoryOutbox(context); endpointConfig.ConfigureSaga<PipelineRunState>(context); }); return configurator; }
生产者配置
private static void AddProducerToBuilder<TMessage>(this IRiderRegistrationConfigurator config, string topicName) where TMessage : class { config.AddProducer<string, TMessage>( topicName, (ctx, cfg) => { cfg.BatchNumMessages = 1; cfg.Linger = TimeSpan.FromMilliseconds(1); cfg.SetValueSerializer(new JsonSerializer<TMessage>()); }); }
状态机事件定义
InstanceState(config => config.CurrentState); Initially( When(FundamentalDataReceived) .Then(context => TransitionToState(context, Initializing)) .Activity(a => a.OfType<InitializationActivity>())); OnUnhandledEvent(async e => await e.Ignore()); DuringAny( When(ProducerValidationResult) .IfElse( ctx => ctx.Message.IsSuccess, successBinder => successBinder .Then(ctx => LogContext.Debug?.Log("[{CorrelationId}]ProducerValidation.Result.Success", ctx.CorrelationId)) .Then(ctx => TransitionToState(ctx, Initialized)) .Activity(a => a.OfInstanceType<RequestMovementActivity>()), failedBinder => failedBinder .Then(ctx => LogContext.Debug?.Log("[{CorrelationId}]ProducerValidation.Result.Fail", ctx.CorrelationId)) .Then(ctx => TransitionToState(ctx, ProducerValidationFailed)) .Activity(a => a.OfInstanceType<RequestQuarantineActivity>())));
流程说明
接收到基础数据后,InitializationActivity向Kafka主题发送消息;对应服务处理完成后返回ProducerValidationResult响应。多数场景下流程正常,但偶发接收响应后状态机事件未触发。
已启用调试日志,查阅官方文档示例后仍未定位问题根源,恳请有相关经验的开发者提供排查思路。
事件未触发的日志记录
状态机日志
| Timestamp | Message |
|---|---|
| 2024-08-08 15:06:20.057 | [41 msec] Published message with CorrelationId: "e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3" |
| 2024-08-08 15:06:20.118 | SAGA:"Test.MyTest.Messaging.StateMachine.Saga.PipelineRunState":e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3 Created "Test.MyTest.Messaging.Messages.Events.IngestFundamental" |
| 2024-08-08 15:06:20.118 | SAGA:"Test.MyTest.Messaging.StateMachine.Saga.PipelineRunState":e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3 Added "Test.MyTest.Messaging.Messages.Events.IngestFundamental" |
| 2024-08-08 15:06:20.118 | [e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3]Transitioning to state: "Initializing (State)" |
| 2024-08-08 15:06:20.118 | [e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3]Requesting producer validation |
| 2024-08-08 15:06:20.167 | SEND loopback://localhost/kafka/test.validation.producer.request 01000000-eec8-261d-ed84-08dcb7b34730 "Test.MyTest.Messaging.Messages.Requests.ProducerValidationRequest" |
| 2024-08-08 15:06:20.167 | [49 msec] Published message with correlationId: e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3 into kafka topic: "test.validation.producer.request" |
| 2024-08-08 15:06:20.227 | RECEIVE loopback://localhost/kafka/test.validation.producer.result 01000000-1f32-a67a-c38e-08dcb7b34738 "Test.MyTest.Messaging.Messages.Events.ProducerValidationResult" "Test.MyTest.Messaging.StateMachine.Saga.PipelineRunState"(00:00:00.0018114) |
生产者验证服务日志
| Timestamp | Message |
|---|---|
| 2024-08-08 15:06:20.167 | Received message with correlationId: "e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3" |
| 2024-08-08 15:06:20.217 | [47 msec] Published message with correlationId: "e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3" into kafka topic: "test.validation.producer.result" |
内容的提问来源于stack exchange,提问作者user26706031
相关产品推荐
相关产品推荐

