如何为 ASP.NET Core 服务提供者编写 EasyNetQ 自动订阅者调度程序?

Ant*_*ean 2 easynetq asp.net-core

所以,EasyNetQ 自动订阅器有一个基本的默认调度器,它不能使用非无参数构造函数创建消息消费者类。

要查看实际效果,请创建一个具有所需依赖项的使用者。您可以设置自己的服务,也可以使用ILogger<T>由框架默认值自动注册的 。

消费文本消息.cs

public class ConsumeTextMessage : IConsume<TextMessage>
{
    private readonly ILogger<ConsumeTextMessage> logger;

    public ConsumeTextMessage(ILogger<ConsumeTextMessage> logger)
    {
        this.logger = logger;
    }

    public void Consume(TextMessage message)
    {
        ...
    }
}
Run Code Online (Sandbox Code Playgroud)

连接自动订阅器(就在哪里/何时编写/运行此代码而言,这里有一些回旋余地)。

启动文件

public void ConfigureServices(IServiceCollection services)
{
    services.AddSingleton<IBus>(RabbitHutch.CreateBus("host=localhost"));
}
Run Code Online (Sandbox Code Playgroud)

(在其他地方,也许startup.Configure或一个BackgroundService)

var subscriber = new AutoSubscriber(bus, "example");
subscriber.Subscribe(Assembly.GetExecutingAssembly());
Run Code Online (Sandbox Code Playgroud)

现在,启动程序并发布一些消息,您应该会看到每条消息都最终出现在默认错误队列中。

System.MissingMethodException: No parameterless constructor defined for this object.
   at System.RuntimeTypeHandle.CreateInstance(RuntimeType type, Boolean publicOnly, Boolean wrapExceptions, Boolean& canBeCached, RuntimeMethodHandleInternal& ctor)
   at System.RuntimeType.CreateInstanceSlow(Boolean publicOnly, Boolean wrapExceptions, Boolean skipCheckThis, Boolean fillCache)
   at System.Activator.CreateInstance[T]()
   at EasyNetQ.AutoSubscribe.DefaultAutoSubscriberMessageDispatcher.DispatchAsync[TMessage,TAsyncConsumer](TMessage message)
   at EasyNetQ.Consumer.HandlerRunner.InvokeUserMessageHandlerInternalAsync(ConsumerExecutionContext context)
Run Code Online (Sandbox Code Playgroud)

我知道我可以提供自己的 dispatcher,但是我们如何与 ASP.NET Core 服务提供者一起使用它;确保这适用于范围服务?

Ant*_*ean 8

所以,这就是我想出的。

public class MessageDispatcher : IAutoSubscriberMessageDispatcher
{
    private readonly IServiceProvider provider;

    public MessageDispatcher(IServiceProvider provider)
    {
        this.provider = provider;
    }

    public void Dispatch<TMessage, TConsumer>(TMessage message)
        where TMessage : class
        where TConsumer : class, IConsume<TMessage>
    {
        using(var scope = provider.CreateScope())
        {
            var consumer = scope.ServiceProvider.GetRequiredService<TConsumer>();
            consumer.Consume(message);
        }
    }

    public async Task DispatchAsync<TMessage, TConsumer>(TMessage message)
        where TMessage : class
        where TConsumer : class, IConsumeAsync<TMessage>
    {
        using(var scope = provider.CreateScope())
        {
            var consumer = scope.ServiceProvider.GetRequiredService<TConsumer>();
            await consumer.ConsumeAsync(message);
        }
    }
}
Run Code Online (Sandbox Code Playgroud)

几个值得注意的点......

该IServiceProvider依赖性是ASP.NET核心DI容器。起初这可能不清楚,因为在整个过程中Startup.ConfigureServices(),您使用另一个接口IServiceCollection.

public MessageDispatcher(IServiceProvider provider)
{
    this.provider = provider;
}
Run Code Online (Sandbox Code Playgroud)

为了解决范围服务,您需要围绕创建和使用使用者创建和管理范围的生命周期。我使用GetRequiredService<T>扩展方法是因为我真的想要一个讨厌的异常,而不是一个可能在我们注意到它之前泄漏一段时间的空引用(以空引用异常的形式)。

using(var scope = provider.CreateScope())
{
    var consumer = scope.ServiceProvider.GetRequiredService<TConsumer>();
    consumer.Consume(message);
}
Run Code Online (Sandbox Code Playgroud)

如果您只provider直接使用 ,如在 中provider.GetRequiredService<T>(),您会在尝试解析范围消费者或消费者的范围依赖时看到这样的错误。

抛出异常:Microsoft.Extensions.DependencyInjection.dll 中的“System.InvalidOperationException”:“无法解析来自根提供程序的范围服务“Example.Messages.ConsumeTextMessage”。

为了解析范围服务并为异步消费者正确维护其生命周期,您需要在正确的位置获取 async/await 关键字。您应该等待ConsumeAsync调用,这要求方法是异步的。在 await 行和您的消费者中使用断点,并逐行地更好地处理这个问题!

public async Task DispatchAsync<TMessage, TConsumer>(TMessage message)
    where TMessage : class
    where TConsumer : class, IConsumeAsync<TMessage>
{
    using(var scope = provider.CreateScope())
    {
        var consumer = scope.ServiceProvider.GetRequiredService<TConsumer>();
        await consumer.ConsumeAsync(message);
    }
}
Run Code Online (Sandbox Code Playgroud)

好的,现在我们有了调度程序,我们只需要在 Startup 中正确设置一切。我们需要从提供者解析调度器,以便提供者可以正确地提供自己。这只是这样做的一种方式。

启动文件

public void ConfigureServices(IServiceCollection services)
{
    // messaging
    services.AddSingleton<IBus>(RabbitHutch.CreateBus("host=localhost"));
    services.AddSingleton<MessageDispatcher>();
    services.AddSingleton<AutoSubscriber>(provider =>
    {
        var subscriber = new AutoSubscriber(provider.GetRequiredService<IBus>(), "example")
        {
            AutoSubscriberMessageDispatcher = provider.GetRequiredService<MessageDispatcher>();
        }
    });

    // message handlers
    services.AddScoped<ConsumeTextMessage>();
}

public void Configure(IApplicationBuilder app, IHostingEnvironment env)
{
    app.ApplicationServices.GetRequiredServices<AutoSubscriber>().SubscribeAsync(Assembly.GetExecutingAssembly());
}
Run Code Online (Sandbox Code Playgroud)