【问题标题】:MassTransit messages types must not be System types exceptionMassTransit 消息类型不能是系统类型异常
【发布时间】:2021-10-27 12:36:17
【问题描述】:

我对 MassTransit 还很陌生,不明白我做错了什么导致以下异常:Messages types must not be System types

这是我的定义:

[BsonIgnoreExtraElements]
public class ArcProcess : SagaStateMachineInstance, ISagaVersion
{
    public Guid CorrelationId { get; set; }

    public string CurrentState { get; set; }

    public int Version { get; set; }

    public Guid ActivationId { get; set; }
}

public static class MessageContracts
{
    static bool _initialized;

    public static void Initialize()
    {
        if (_initialized)
            return;

        GlobalTopology.Send.UseCorrelationId<StartProcessingMessage>(x => x.ActivationId);
        GlobalTopology.Send.UseCorrelationId<ReconstructionFinishedMessage>(x => x.ActivationId);
        GlobalTopology.Send.UseCorrelationId<ProcessingFinishedMessage>(x => x.ActivationId);

        _initialized = true;
    }
}

我的 2 个消费者是:

public class StartReconstructionConsumer : IConsumer<StartProcessingMessage>
{
    readonly ILogger<StartReconstructionConsumer> _Logger;

    private readonly int _DelaySeconds = 5;

    public StartReconstructionConsumer(ILogger<StartReconstructionConsumer> logger)
    {
        _Logger = logger;
    }

    public async Task Consume(ConsumeContext<StartProcessingMessage> context)
    {
        var activationId = context.Message.ActivationId;

        _Logger.LogInformation($"Received Scan: {activationId}");

        await Task.Delay(_DelaySeconds * 1000);

        _Logger.LogInformation($"Finish Scan: {activationId}");

        await context.Publish<ReconstructionFinishedMessage>(new { ActivationId = activationId });
    }
}

public class ProcessingFinishedConsumer : IConsumer<ProcessingFinishedMessage>
{
    readonly ILogger<ProcessingFinishedConsumer> _Logger;

    public ProcessingFinishedConsumer(ILogger<ProcessingFinishedConsumer> logger)
    {
        _Logger = logger;
    }

    public async Task Consume(ConsumeContext<ProcessingFinishedMessage> context)
    {
        _Logger.LogInformation($"Finish {context.Message.ActivationId}");

        await Task.CompletedTask;
    }
}

这是 StateMachine 的定义:

public class ArcStateMachine: MassTransitStateMachine<ArcProcess>
{
    static ArcStateMachine()
    {
        MessageContracts.Initialize();
    }

    public ArcStateMachine()
    {
        InstanceState(x => x.CurrentState);

        Initially(
            When(ProcessingStartedEvent)
            .Then(context =>
            {
                Console.WriteLine(">> ProcessingStartedEvent");
                context.Instance.ActivationId = context.Data.ActivationId;
            })
            .TransitionTo(ProcessingStartedState));

        During(ProcessingStartedState,
            When(ReconstructionFinishedEvent)
            .Then(context =>
            {
                Console.WriteLine(">> ReconstructionFinishedEvent");
                context.Instance.ActivationId = context.Data.ActivationId;
            })
            .Publish(context =>
            {
                return context.Init<ProcessingFinishedMessage>(new { ActivationId = context.Data.ActivationId });
            })
            .TransitionTo(ProcessingFinishedState)
            .Finalize());
    }

    public State ProcessingStartedState { get; }
    public State ReconstructionStartedState { get; }
    public State ReconstructionFinishedState { get; }
    public State ProcessingFinishedState { get; }

    public Event<StartProcessingMessage> ProcessingStartedEvent { get; }
    public Event<ReconstructionStartedMessage> ReconstructionStartedEvent { get; }
    public Event<ReconstructionFinishedMessage> ReconstructionFinishedEvent { get; }
    public Event<ProcessingFinishedMessage> ProcessingFinishedEvent { get; }
}

MassTransit 的设置如下所示:

var rabbitHost = Configuration["RABBIT_MQ_HOST"];

        if (rabbitHost.IsNotEmpty())
        {
            services.AddMassTransit(cnf =>
            {
                var connectionString = Configuration["MONGO_DB_CONNECTION_STRING"];

                var machine = new ArcStateMachine();
                var repository = MongoDbSagaRepository<ArcProcess>.Create(connectionString,
                    "mongoRepo", "WorkflowState");

                cnf.AddConsumer(typeof(StartReconstructionConsumer));
                cnf.AddConsumer(typeof(ProcessingFinishedConsumer));

                cnf.UsingRabbitMq((context, cfg) =>
                {
                    cfg.Host(new Uri(rabbitHost), hst =>
                    {
                        hst.Username("guest");
                        hst.Password("guest");
                    });

                    cfg.ConfigureEndpoints(context);

                    cfg.ReceiveEndpoint(BusConstants.SagaQueue,
                        e => e.StateMachineSaga(machine, repository));
                });
            });

            services.AddMassTransitHostedService();

            services.AddSwaggerGen(c =>
            {
                c.SwaggerDoc("v1", new OpenApiInfo { Title = "MyApp", Version = "v1" });
            });
        }

我有几个问题:

  1. 实际上何时发布消息作为发布消息的结果? IE。在我的示例中,await _BusInstance.Bus.Publish&lt;StartProcessingMessage&gt;(new { ActivationId = id }); 是从 StartReconstructionConsumer 使用的 WebApi 调用的,但实际上当状态机开始与 Initially(When(ProcessingStartedEvent)... 一起操作时?

  2. 我的处理应确保我已经处于ProcessingStartedState 状态,以便During(ProcessingStartedState, When(ReconstructionFinishedEvent)... 正确行事。那么,如何确保在收到StartProcessingMessage 时触发的消费者可以发布应该启动DuringReconstructionFinishedMessage?我是否正确构建了消息交换?

  3. 目前对于await context.Publish&lt;ReconstructionFinishedMessage&gt;(new { ActivationId = activationId });,我在日志中发现一个异常,指出R-FAULT rabbitmq://localhost/saga.service d4070000-7b3b-704d-0f10-08d99942c959 Nanox.GC.Shared.AppCore.Messages.ReconstructionFinishedMessage ReconCaller.Saga.ArcProcess(00:00:04.1132604),而消息中的guid实际上是MessageId。我在rabbitmq中的消息被路由到saga.service_error,但Messages types must not be System types: System.Threading.Tasks.Task&lt;Nanox.GC.Shared.AppCore.Messages.ProcessingFinishedMessage&gt; (Parameter 'T')除外。

看来我在这里真的很想念..

我的意图是启动处理,该处理将由几个消费者按顺序处理几个阶段。所以在这里我尝试构建一个简单的状态机,只要有人打电话给StartProcessing,它就会启动,然后每个消费者都会完成它的工作并触发FinishedStepX,这会将状态机提升到一个新的步骤并启动下一个消费者,直到所有处理完成,状态机将报告ProcessingComplete

感谢您的任何帮助和进步

【问题讨论】:

    标签: c# masstransit


    【解决方案1】:

    首先,你的总线配置有点奇怪,所以我已经清理了:

    services.AddMassTransit(cnf =>
    {
        var connectionString = Configuration["MONGO_DB_CONNECTION_STRING"];
    
        cfg.AddSagaStateMachine<ArcStateMachine, ArcProcess>()
            .Endpoint(e => e.Name = BusConstants.SagaQueue)
            .MongoDbRepository(connectionString, r =>
            {
                r.DatabaseName = "mongoRepo";
                r.CollectionName = "WorkflowState";
            });
    
        cnf.AddConsumer<StartReconstructionConsumer>();
        cnf.AddConsumer<ProcessingFinishedConsumer>();
    
        cnf.UsingRabbitMq((context, cfg) =>
        {
            cfg.Host(new Uri(rabbitHost), hst =>
            {
                hst.Username("guest");
                hst.Password("guest");
            });
    
            cfg.ConfigureEndpoints(context);
        });
    });
    

    并且发布问题与使用的方法有关,只有 PublishAsync 允许使用消息初始化器:

    During(ProcessingStartedState,
        When(ReconstructionFinishedEvent)
            .Then(context =>
            {
                Console.WriteLine(">> ReconstructionFinishedEvent");
                context.Instance.ActivationId = context.Data.ActivationId;
            })
            .PublishAsync(context =>
            {
                return context.Init<ProcessingFinishedMessage>(new { ActivationId = context.Data.ActivationId });
            })
            .TransitionTo(ProcessingFinishedState)
            .Finalize());
    

    这应该可以解决您的问题。

    【讨论】:

    • 感谢您的快速回复(小修复:cnf.AddSaga..)。修复了问题。现在有一些新的东西:发布ProcessingFinishedMessage 后,我在日志中收到:R-FAULT rabbitmq://localhost/saga.service d4070000-7b3b-704d-e7df-08d9994b2062 Nanox.GC.Shared.AppCore.Messages.ProcessingFinishedMessage ReconCaller.Saga.ArcProcess(00:00:00.0286223),消息转到saga.service_error,并带有以下错误消息:The ProcessingFinishedEvent event is not handled during the ProcessingStartedState state for the ArcStateMachine state machine。我应该为此添加一些.Ignore() 吗?
    • 你为什么要在你的状态机中处理已完成的处理?不就是制造事件吗?从状态机中删除它,您可能必须删除 RMQ 中的交换绑定。
    • 那么系统如何知道处理已经完成了所有阶段?这就是为什么我有ProcessingFinishedConsumer 应该更新处理已成功通过所有阶段的系统/用户(无论如何)。而且我还希望状态机将状态更新为ProcessingFinishedState,以便将其反映在消息存储库中。我错过了什么吗?
    • 好吧,调用Finalize 转换为Final。如果您希望状态机保留ProcessingFinishedState,请不要调用Finalize。没有理由使用事件来改变状态,它正在产生事件——它已经知道了。
    • 关于致电Finalize 我了解。因此,或者我可以摆脱ProcessingFinishedState,所以我只是删除了.TransitionTo(ProcessingFinishedState),只使用.Finalize());,但我仍然在RMQ中遇到同样的错误,因为(我想)发布ProcessingFinishedMessage [之前finalization] 由ProcessingFinishedConsumer 使用,我确实需要它,以便我的系统了解状态机何时完成。我应该以不同的方式实现它吗?
    【解决方案2】:

    在@Chris Patterson 的慷慨帮助下,可行的解决方案是:

    定义:

    [BsonIgnoreExtraElements]
    public class ArcProcess : SagaStateMachineInstance, ISagaVersion
    {
        public Guid CorrelationId { get; set; }
    
        public string CurrentState { get; set; }
    
        public int Version { get; set; }
    
        public Guid ActivationId { get; set; }
    }
    
    public interface StartProcessingMessage
    {
        Guid ActivationId { get; }
    }
    
    public interface ProcessingFinishedMessage
    {
        Guid ActivationId { get; }
    }
    
    public static class MessageContracts
    {
        static bool _initialized;
    
        public static void Initialize()
        {
            if (_initialized)
                return;
    
            GlobalTopology.Send.UseCorrelationId<StartProcessingMessage>(x => x.ActivationId);
            GlobalTopology.Send.UseCorrelationId<ProcessingFinishedMessage>(x => x.ActivationId);
    
            _initialized = true;
        }
    }
    

    消费者:

    public class StartProcessingConsumer : IConsumer<StartProcessingMessage>
    {
        readonly ILogger<StartProcessingConsumer> _Logger;
    
        private readonly int _DelaySeconds = 5;
    
        public StartProcessingConsumer(ILogger<StartProcessingConsumer> logger)
        {
            _Logger = logger;
        }
    
        public async Task Consume(ConsumeContext<StartProcessingMessage> context)
        {
            var activationId = context.Message.ActivationId;
    
            _Logger.LogInformation($"Received Scan: {activationId}");
    
            await Task.Delay(_DelaySeconds * 1000);
    
            _Logger.LogInformation($"Finish Scan: {activationId}");
    
            await context.Publish<ProcessingFinishedMessage>(new { ActivationId = activationId });
        }
    }
    
    public class ProcessingFinishedConsumer : IConsumer<ProcessingFinishedMessage>
    {
        readonly ILogger<ProcessingFinishedConsumer> _Logger;
    
        public ProcessingFinishedConsumer(ILogger<ProcessingFinishedConsumer> logger)
        {
            _Logger = logger;
        }
    
        public async Task Consume(ConsumeContext<ProcessingFinishedMessage> context)
        {
            _Logger.LogInformation($"Finish {context.Message.ActivationId}");
    
            await Task.CompletedTask;
        }
    }
    

    状态机定义:

    public class ArcStateMachine: MassTransitStateMachine<ArcProcess>
    {
        static ArcStateMachine()
        {
            MessageContracts.Initialize();
        }
    
        public ArcStateMachine()
        {
            InstanceState(x => x.CurrentState);
    
            Initially(
                When(ProcessingStartedEvent)
                .Then(context =>
                {
                    context.Instance.ActivationId = context.Data.ActivationId;
                })
                .TransitionTo(ProcessingStartedState));
    
            During(ProcessingStartedState,
                When(ProcessingFinishedEvent)
                .Then(context =>
                {
                    context.Instance.ActivationId = context.Data.ActivationId;
                })
                .Finalize());
        }
    
        public State ProcessingStartedState { get; }
        public State ProcessingFinishedState { get; }
    
        public Event<StartProcessingMessage> ProcessingStartedEvent { get; }
        public Event<ProcessingFinishedMessage> ProcessingFinishedEvent { get; }
    }
    

    MassTransit 设置:

    var rabbitHost = Configuration["RABBIT_MQ_HOST"];
    
            if (rabbitHost.IsNotEmpty())
            {
                services.AddMassTransit(cnf =>
                {
                    var connectionString = Configuration["MONGO_DB_CONNECTION_STRING"];
    
                    cnf.AddSagaStateMachine<ArcStateMachine, ArcProcess>()
                        .Endpoint(e => e.Name = BusConstants.SagaQueue)
                        .MongoDbRepository(connectionString, r =>
                        {
                            r.DatabaseName = "mongoRepo";
                            r.CollectionName = "WorkflowState";
                        });
    
    
                    cnf.AddConsumer(typeof(StartProcessingConsumer));
                    cnf.AddConsumer(typeof(ProcessingFinishedConsumer));
    
                    cnf.UsingRabbitMq((context, cfg) =>
                    {
                        cfg.Host(new Uri(rabbitHost), hst =>
                        {
                            hst.Username("guest");
                            hst.Password("guest");
                        });
    
                        cfg.ConfigureEndpoints(context);
                    });
                });
    
                services.AddMassTransitHostedService();
    
                services.AddSwaggerGen(c =>
                {
                    c.SwaggerDoc("v1", new OpenApiInfo { Title = "MyApp", Version = "v1" });
                });
            }
    

    这个例子对我理解 MassTrasit 的基本原理有很大帮助。

    【讨论】:

    • 感谢您展示您的最终作品。所以现在一切都很好?
    • 是的,非常完美。我还添加了一个中间阶段来检查 2 个阶段的执行情况,它工作正常。我当前的任务是找到一种方法来处理消费者可能发生的任何异常。我收到的结果是创建了 [queue_name]_error 并且消息出现在那里。我添加了一个 Fault 消费者,但现在我有两个问题:a) 我怎样才能让这个消费者处理任何消息类型的任何故障? b) 让错误消费者处理的消息从_error队列中消失?
    • 失败的消费者发布的Fault&lt;T&gt; 与移动到_error 队列的原始消息是分开的。错误队列的管理是一个单独的问题,不由 MassTransit 处理。
    • 感谢您的解释。我的问题是是否有办法让Fault&lt;&gt; 消费者处理消费消息的代码中的任何类型的故障/异常,但防止消息被路由到_error 队列?
    • 您可以在接收端点配置上指定DiscardFaultedMessages(),但请谨慎操作,因为错误消息将丢失。
    猜你喜欢
    • 2018-10-29
    • 2016-10-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多