MassTransit配置疑问:如何避免Kafka消费的消息被推送至RabbitMQ
首先得帮你澄清一个误解:你日志里看到的rabbitmq://my-rabbitmq/kafka/my-kafka-topic/my-service并不是真的把Kafka消息推送到了RabbitMQ,这只是MassTransit的逻辑端点命名规则导致的——因为你把Kafka Rider挂载到了以RabbitMQ为主总线的实例上,所以MassTransit会用主总线的地址格式来命名Kafka消费者的端点,但实际上Kafka的消息并不会被转发到RabbitMQ集群里,你发现没有对应的Exchange被创建,也能佐证这一点。
不过如果你想彻底消除这种关联、避免后续混淆,完全可以通过配置把RabbitMQ和Kafka的消费逻辑隔离开,这里给你两种靠谱的方案:
方案一:使用两个独立的MassTransit总线实例(推荐)
最彻底的方式是分别为RabbitMQ和Kafka创建独立的MassTransit总线,让两者完全没有关联。这样Kafka的消费逻辑完全属于自己的总线,不会和RabbitMQ有任何绑定。
配置RabbitMQ总线
if (options.RabbitMq.Enabled) { services.AddMassTransit(busConfigurator => { busConfigurator.SetKebabCaseEndpointNameFormatter(); busConfigurator.AddConsumer<RabbitMqConsumer, RabbitMqConsumerDefinition>(); busConfigurator.UsingRabbitMq((context, busFactoryConfigurator) => { busFactoryConfigurator.Host(new Uri(options.RabbitMq.Endpoint)); busFactoryConfigurator.ConfigureEndpoints(context); }); }); }
配置Kafka专属总线(用InMemory作为主传输)
因为你的Kafka消费者不需要生产消息,所以用InMemory作为Kafka总线的主传输就足够了,这样完全不会和RabbitMQ产生关联:
if (options.Kafka.Enabled) { var schemaRegistryConfig = new SchemaRegistryConfig { Url = options.Kafka.SchemaRegistryUrl }; var schemaRegistryClient = new CachedSchemaRegistryClient(schemaRegistryConfig); services.AddMassTransit(busConfigurator => { busConfigurator.AddSingleton<ISchemaRegistryClient>(schemaRegistryClient); busConfigurator.AddRider(rider => { rider.AddConsumer<KafkaMyMessageConsumer>(); rider.UsingKafka((context, kafkaConfigurator) => { kafkaConfigurator.ClientId = options.Kafka.ClientId; kafkaConfigurator.Host(options.Kafka.Endpoint, configureHost => { // 保留你的ssl、sasl配置 }); kafkaConfigurator.TopicEndpoint<string, MyMessage>(options.Kafka.Topic, options.Kafka.GroupId, endpointConfigurator => { endpointConfigurator.CreateIfMissing(t => t.NumPartitions = 1); endpointConfigurator.SetValueDeserializer(new ProtobufDeserializer<MyMessage>().AsSyncOverAsync()); endpointConfigurator.AutoOffsetReset = AutoOffsetReset.Earliest; endpointConfigurator.ConfigureConsumer<KafkaMyMessageConsumer>(context); }); }); }); // 用InMemory作为Kafka总线的主传输,彻底隔离RabbitMQ busConfigurator.UsingInMemory((context, cfg) => { cfg.ConfigureEndpoints(context); }); }); }
方案二:调整Rider的端点配置(适合不想拆分总线的场景)
如果你不想创建两个总线,也可以修改Kafka Rider的端点配置,不让它使用主总线(RabbitMQ)的命名规则:
// 在Kafka Rider的配置里添加这一行,自定义端点命名,避免关联RabbitMQ rider.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("kafka", false));
不过这种方式还是会依赖主总线,不如第一种方案彻底。
你之前尝试的“用两次AddMassTransit+InMemory总线”的思路是完全正确的,日志里出现的loopback uri就是InMemory总线的端点,这是正常现象,而且能彻底避免Kafka和RabbitMQ的关联问题。
备注:内容来源于stack exchange,提问作者esskar

