【问题标题】:Add Consumers without adding MassTransit to a dependency injection container添加消费者而不将 MassTransit 添加到依赖注入容器
【发布时间】:2020-06-26 10:48:32
【问题描述】:

我正在尝试实现与传输类型(即 Azure 服务总线、RabbitMQ 等)无关的 MassTransit。我已经将逻辑与 DI 容器分离,因为 DI 容器将接口的范围添加到具体类,并具有一个工厂来根据配置确定传输。这适用于发布和发送,但是当我尝试执行请求/响应或消费时,会创建 RabbitMq 的交换/队列以及 Azure 服务总线的主题和队列,但没有收到任何消息。

我的想法是,然后我可以使用在项目之间执行所有逻辑和端口的类库,而不需要所有消费者的大量 sn-p,因为这些消费者可能因项目而异。

我觉得好像我错过了一些重要的东西。在使用 AddMassTransit 将总线配置添加到 DI 容器中之后,如果不将 .AddConsumer 放入 DI 容器中,这将不起作用,这与我正在尝试的相反。

我构建了一个集成测试项目来尝试 TDD,同时在添加到服务之前解决此功能并在此处被阻止。响应变量永远不会返回,我最终会超时。集成项目仅依赖于 MassTransit 类库和 Configuration 类库,实际上是直接与传输通信(而不是模拟它)。

我是接近了还是需要放弃这个任务?

我当前项目的代码方面:

Startup.cs

        Configuration.IConfigurationProvider cfg = new ServiceFabricConfigurationProvider();
        switch (cfg.MessageService.MessagingService)
        {
            case "AzureServiceBus":
                services.AddScoped<IMassTransitTransport, MassTransitAzureServiceBusTransport>();
                break;
            case "RabbitMq":
                services.AddScoped<IMassTransitTransport, MassTransitRabbitMqTransport>();
                break;
            default:
                throw new ArgumentException("Invalid message service");
        };

        services.AddScoped<IMessagingService, MassTransitMessagingService>();

Transport 类非常相似,因为它们有一个共同的接口 - 如下所示

IMassTransitTransport

public interface IMassTransitTransport
{
    IBusControl BusControl { get; }
}

MassTransitAzureServiceBusTransport

public sealed class MassTransitAzureServiceBusTransport : IMassTransitTransport
{
    readonly IConfigurationProvider configProvider;

    public MassTransitAzureServiceBusTransport(IConfigurationProvider configProvider)
    {
        this.configProvider = configProvider;
        BusControl = ConfigureBus();
        BusControl.StartAsync();
    }

    public IBusControl BusControl { get; }

    IBusControl ConfigureBus()
    {
        return Bus.Factory.CreateUsingAzureServiceBus(cfg => 
        {
            cfg.Host(configProvider.AzureServiceBus.AzureServiceBusConnectionString);

            cfg.ReceiveEndpoint("MyQueue", e => 
            {
                e.Consumer<ConsumerClass>();
            });
        });
    }

MassTransitRabbitMqTransport

public sealed class MassTransitRabbitMqTransport : IMassTransitTransport
{
    readonly IConfigurationProvider configProvider;

    public MassTransitRabbitMqTransport(IConfigurationProvider configProvider)
    {
        this.configProvider = configProvider;
        BusControl = ConfigureBus();
        BusControl.StartAsync();
    }

    public IBusControl BusControl { get; }

    IBusControl ConfigureBus()
    {
        return Bus.Factory.CreateUsingRabbitMq(cfg =>
        {

            cfg.Host(new Uri(configProvider.Rabbit.HostAddress), host =>
            {
                host.Username(configProvider.Rabbit.Username);
                host.Password(configProvider.Rabbit.Password);
            });

            cfg.ReceiveEndpoint("MyQueue", e => 
            {
                e.Consumer<ConsumerClass>();
            });
        });
    }
}

消息服务

public interface IMessagingService
{
    Task Publish<T>(object payload) where T : class;
    Task Send<T>(object payload) where T : class;
}

public class MassTransitMessagingService : IMessagingService
{
    readonly IMassTransitTransport massTransitTransport;

    public MassTransitMessagingService(IMassTransitTransport massTransitTransport)
    {
        //transport bus config already happens in massTransitTransport constructor
        this.massTransitTransport = massTransitTransport;
    }

    public async Task Publish<T>(object payload) where T : class
    {
        await massTransitTransport.BusControl.Publish<T>(payload);
    }

    public async Task Send<T>(object payload) where T : class
    {
        var endpoint = await massTransitTransport.BusControl.GetSendEndpoint(new Uri(massTransitTransport.BusControl.Address, typeof(T).ToString()));
        await endpoint.Send<T>(payload);
    }
}

请求/响应接口和类

public interface IEventRequest
{
    Guid EventGuid { get; set; }
    string Message { get; set; }
}

public interface IEventResponse
{
    Guid EventGuid { get; set; }
    string RequestMessage { get; set; }
    string ResponseMessage { get; set; }
}

public class ConsumerClass : IConsumer<IEventResponse>
{
    public async Task Consume(ConsumeContext<IEventResponse> context)
    {
        var payloadResponse = new
        {
            context.Message.EventGuid,
            context.Message.RequestMessage,
            ResponseMessage = "This is the response message;"
        };

        await context.RespondAsync(payloadResponse);
    }
}

执行请求/响应的测试方法

    [Test]
    public async Task SendAndReceiveMessage()
    {
        // arrange
        var config = GetConfiguration();
        var transport = new MassTransitRabbitMqTransport(config);

        var payload = new
        {
            EventGuid = Guid.NewGuid(),
            Message = "This is an event message"
        };

        var clientFactory = transport.BusControl.CreateClientFactory();
        var client = clientFactory.CreateRequestClient<IEventRequest>();
        var response = await client.GetResponse<ConsumerClass>(payload);

    }

更新 1

根据 Chris Patterson 的反馈,我进行了以下修改。我相信我通过显示启动类的片段引起了混乱。

这里实际上发生了两件事:一个 API 项目,其中包含上面的部分和下面的重构代码片段,以及一个集成测试项目。

在集成测试项目方面,确实只用到了3个类库。集成测试项目不包含 DI 容器,因为它是一个测试项目,我想看看是否可以将逻辑与 DI 容器解耦。另外,集成测试项目中没有 ILogger。

整个解决方案是一个 API 项目,它确实有一个带有内置 DI 容器的启动和一个实现发布的控制器。我的尝试是将 MassTransit 逻辑与 DI 容器分离,以便可以在其他地方使用 MassTransitTransport 项目,从而可以使用 MassTransit 支持的任何传输。我的问题是,在使用消息方面这是否是一个坏主意(即除非使用 DI 容器,否则无法完成)或者是否可以完成。如果可以做到,我错过了什么/我有什么问题?

配置 - 利用 IConfigurationProvider 大众运输 - 包含 ConsumerClass、IEventRequest、IEventResponse、EventResponse、 IMassTransitTransport、MassTransitTransportRabbitMqTransport、 MassTransitZaureServiceBusTransport、IMessagingService、 MassTransitMessagingService MassTransitTransport.Integration - 包含 MassTransitMessagingServiceTests, NonServiceFabricConfigurationProvider

创建 EventResponse 并在 ConsumerClass 中使用

public class EventResponse : IEventResponse
{
    public Guid EventGuid { get; set; }
    public string RequestMessage { get; set; }
    public string ResponseMessage { get; set; }
}

消费类

public class ConsumerClass : IConsumer<IEventResponse>
{
    public async Task Consume(ConsumeContext<IEventResponse> context)
    {
        var payloadResponse = new EventResponse()
        {
            EventGuid = context.Message.EventGuid,
            RequestMessage = context.Message.RequestMessage,
            ResponseMessage = "This is the response message;"
        };

        await context.RespondAsync(payloadResponse);
    }
}

集成测试项目

    [Test]
    public async Task SendAndReceiveMessage()
    {
        // arrange
        var config = GetConfiguration();
        var transport = new MassTransitRabbitMqTransport(config);
        //var transport = new MassTransitAzureServiceBusTransport(config);

        var payload = new
        {
            EventGuid = Guid.NewGuid(),
            Message = "This is an event message"
        };

        var clientFactory = transport.BusControl.CreateClientFactory();
        var client = clientFactory.CreateRequestClient<IEventRequest>();
        var response = await client.GetResponse<IEventResponse>(payload);
    }

收到错误消息

MassTransit.RequestTimeoutException : 等待响应超时,RequestId: 7a000000-9a3c-0005-8037-08d80882f498

【问题讨论】:

  • 您的payloadResponse 是匿名类型,这是不允许的。启用日志记录(通过 ILoggerFactory,请参阅文档)会向您显示错误。您需要使用 ResponseAsync().
  • 另外,GetResponse&lt;IEventResponse&gt;() - 消费者不是响应的消息类型。
  • 另外,我会避免启动总线,或者在服务配置中一般做任何 IO。 ConfigureServices 方法是配置 东西,而不是启动基础设施。您可以注册托管服务以启动和停止总线services.AddSingleton&lt;IHostedService&gt;(new BusHostedService(bus))
  • @ChrisPatterson 我已经进行了一些重构,但仍然出现超时。我添加了一个更新 1 试图澄清我的问题,因为我有误导性。我试图将 MassTransit 项目与 DI 容器尽可能地分离,因为我的计划是将该组件插入其他解决方案。如果这不是一个好主意,请告诉我(即 MassTransit 取决于 DI 容器?)
  • 我建议使用测试工具进行单元测试,如documentation 中所述。您将太多不相关的问题拖入测试中。并且通过尝试不可知论来增加自己周围的复杂性 - 也许考虑在仅依赖于 MassTransit 的业务逻辑之间划定界限,然后为您的运输集成提供单独的程序集。

标签: c# rabbitmq azureservicebus masstransit


【解决方案1】:

我将在评论中添加答案以关闭此问题:

有 2 个问题。第一个是 StartAsync 而不是 Start。第二个是 ConsumerClass 没有正确实现。它应该从 IConsumer 继承,并将方法作为公共 async Task Consume(ConsumeContext context)。之后,一切正常。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-07-22
    相关资源
    最近更新 更多