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