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

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响应。多数场景下流程正常,但偶发接收响应后状态机事件未触发。

已启用调试日志,查阅官方文档示例后仍未定位问题根源,恳请有相关经验的开发者提供排查思路。

事件未触发的日志记录

状态机日志

TimestampMessage
2024-08-08 15:06:20.057[41 msec] Published message with CorrelationId: "e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3"
2024-08-08 15:06:20.118SAGA:"Test.MyTest.Messaging.StateMachine.Saga.PipelineRunState":e1fa4086-8fc1-4d1d-b4c1-8dab3a4bd0f3 Created "Test.MyTest.Messaging.Messages.Events.IngestFundamental"
2024-08-08 15:06:20.118SAGA:"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.167SEND 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.227RECEIVE 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)

生产者验证服务日志

TimestampMessage
2024-08-08 15:06:20.167Received 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 00:01:04