【问题标题】:Masstransit add consumer in container by specifying queue bindingMasstransit 通过指定队列绑定在容器中添加消费者
【发布时间】:2020-03-09 16:05:03
【问题描述】:

我有以下消费者类型:

internal class ObjectAddedHandler : IConsumer<ObjectAddedIntegrationEvent>
{
    public async Task Consume(ConsumeContext<ObjectAddedIntegrationEvent> context)
    {
        var @event = context.Message;
        await HandleAsync(@event).ConfigureAwait(false);
    }
}

通过以下方式在我的容器中注册:

container.AddMassTransit(x =>
{
    x.AddBus(context => Bus.Factory.CreateUsingRabbitMq(cfg =>
    {
        var host = cfg.Host(configurationProvider.RabbitMQHostName, hostConfigurator =>
        {
            hostConfigurator.Username(configurationProvider.RabbitMQUsername);
            hostConfigurator.Password(configurationProvider.RabbitMQPassword);

            hostConfigurator.UseCluster(c =>
            {
                string[] hostnames = configurationProvider.RabbitMQNodes.Split(';');
                c.ClusterMembers = hostnames;
            });
        });

        host.Settings.GetConnectionFactory().Endpoint.AddressFamily = AddressFamily.InterNetwork;

        /*HERE*/ x.AddConsumer<ObjectAddedHandler>().Endpoint(e => e.Name = "ObjectAddedHandler "+configurationProvider.TenantName);

        cfg.ConfigureEndpoints(container);
    }));

});

但是,按照documentation,我想设置一个直接交换,以便使用路由键。在文档中的任何地方都找不到如何以我的方式添加消费者,同时按照文档中的报告设置端点的绑定属性。

当我尝试在添加消费者时访问端点时,我只能修改名称、预取计数和其他几个属性,但仅此而已。但是,我想将我的端点设置为仅接受带有路由键 tenantName 的消息。有什么办法吗?

编辑

发布方:

container.AddMassTransit(x =>
{
    x.AddBus(() => Bus.Factory.CreateUsingRabbitMq(cfg =>
    {

        cfg.Send<ObjectAddedIntegrationEvent>(routingCfg => {
            routingCfg.UseRoutingKeyFormatter(config => ConfigurationValuesProvider.Current.Get("TenantCode"));
        });
        cfg.Message<ObjectAddedIntegrationEvent>(routingCfg => routingCfg.SetEntityName("ObjectAddedIntegrationEvent"));
        cfg.Publish<ObjectAddedIntegrationEvent>(routingCfg => routingCfg.ExchangeType = ExchangeType.Direct);


        var host = cfg.Host(ConfigurationValuesProvider.Current.Get("RabbitMQHostName"), hostConfigurator =>
        {
#if !DEBUG
            hostConfigurator.Username(ConfigurationValuesProvider.Current.Get("RabbitMQUsername"));
            hostConfigurator.Password(ConfigurationValuesProvider.Current.Get("RabbitMQPassword"));

            hostConfigurator.UseCluster(c =>
            {
                string[] hostnames = ConfigurationValuesProvider.Current.Get("RabbitMQNodes").Split(';');
                c.ClusterMembers = hostnames;
            });
#endif
        });
    }));
});

接收方:

container.AddMassTransit(x =>
{
    x.AddBus(context => Bus.Factory.CreateUsingRabbitMq(cfg =>
    {
        var host = cfg.Host(configurationProvider.RabbitMQHostName, hostConfigurator =>
        {
            hostConfigurator.Username(configurationProvider.RabbitMQUsername);
            hostConfigurator.Password(configurationProvider.RabbitMQPassword);

            hostConfigurator.UseCluster(c =>
            {
                string[] hostnames = configurationProvider.RabbitMQNodes.Split(';');
                c.ClusterMembers = hostnames;
            });
        });

        host.Settings.GetConnectionFactory().Endpoint.AddressFamily = AddressFamily.InterNetwork;

        x.AddConsumer<ObjectAddedHandler>().Endpoint(e => e.Name = "ObjectAddedHandler "+configurationProvider.TenantName);

        cfg.ConfigureEndpoints(container);
    }));

});



internal class ObjectAddedHandler : IConsumer<ObjectAddedIntegrationEvent>
{
    public async Task Consume(ConsumeContext<ObjectAddedIntegrationEvent> context)
    {
        var @event = context.Message;
        await HandleAsync(@event).ConfigureAwait(false);
    }
}


internal class ObjectAddedHandlerConsumerDefinition :
            ConsumerDefinition<ObjectAddedHandler>
{
    private readonly IConfigurationProvider _provider;

    public ObjectAddedHandlerConsumerDefinition(IConfigurationProvider provider)
    {
        _provider = provider;

        EndpointName = "ObjectAddedHandler" + provider.TenantName;
    }

    protected override void ConfigureConsumer(IReceiveEndpointConfigurator endpointConfigurator,
    IConsumerConfigurator<ObjectAddedHandler> consumerConfigurator)
    {
        if (endpointConfigurator is IRabbitMqReceiveEndpointConfigurator rabbit)
        {
            rabbit.BindMessageExchanges = false;

            rabbit.Bind("ObjectAddedIntegrationEvent", s =>
            {
                s.RoutingKey = _provider.TenantName;
                s.ExchangeType = ExchangeType.Direct;
            });
        }
    }
}

当消费者启动时,ObjectAddedIntegrationEvent 直接交换被正确创建。 ObjectAddedHandlerTenant 扇出交换(应该是匹配交换)也被正确创建。 不幸的是,当我尝试从发布方发送消息并监控 ObjectAddedIntegrationEvent 直接交换时,我看不到任何消息。

【问题讨论】:

    标签: .net rabbitmq masstransit


    【解决方案1】:

    您可以使用消费者定义来做到这一点,该定义应使用AddConsumer&lt;T&gt;(typeof(definitionclass)) 与您的消费者一起添加。可以类似这样:

    public class ObjectAddedHandlerDefinition :
        ConsumerDefinition<ObjectAddedHandler>
    {
        public ObjectAddedHandlerDefinition(IConfigurationProvider provider)
        {
            _provider = provider;
    
            EndpointName = "ObjectAddedHandler" + provider.TenantName;
    
            ConcurrentMessageLimit = 4;
        }
    
        protected override void ConfigureConsumer(IReceiveEndpointConfigurator endpointConfigurator,
            IConsumerConfigurator<ObjectAddedHandler> consumerConfigurator)
        {
            if(endpointConfigurator is IRabbitMqReceiveEndpointConfigurator rabbit)
            {
                rabbit.BindMessageExchanges = false;
    
                // or use Bind<T> for message type name
                rabbit.Bind("some-exchange", s => 
                {
                    s.RoutingKey = _provider.TenantName;
                    s.ExchangeType = ExchangeType.Direct;
                });
            }
        }
    }
    

    【讨论】:

    • 首先感谢,因为这已经有所帮助。不幸的是,它还没有工作。我也向出版商更新了我的问题。
    • 你去,是的,你需要使用这种方法发布。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-07-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多