MassTransit+RabbitMQ(.NET6):无信封JSON消息如何避免触发所有消费者
我用.NET6结合MassTransit实现RabbitMQ消息功能,代码里把ProductConsumer和NameConsumer配置到同一个接收队列。但从RabbitMQ管理后台直接发消息时,不管payload格式匹配哪个消费者,两个消费者都会被触发。日志显示ProductMessage格式的消息同时触发了NameConsumer(Name字段被设为null)。不想用MassTransit的信封包装(里面冗余信息太多),有没有其他方案?比如通过消息头区分目标消费者?
日志信息:
2024-08-27 11:47:50.3704 INFO [ProductConsumer] Message received context: {"Message":{"Code":"_code","Name":"_name"}} 2024-08-27 11:47:50.4210 INFO [NameConsumer] Message received context: {"Message":{"Name":null}}
原代码示例:
Program.cs
services.AddMassTransit(x => { x.AddConsumer<ProductConsumer>(); x.AddConsumer<NameConsumer>(); x.UsingRabbitMq((context, cfg) => { cfg.Host($"myRabbitHost.com", 5672, "VHOST" , configurator => { configurator.Username("_user"); configurator.Password("_pwd"); }); cfg.ReceiveEndpoint("AUTO.QUEUENAME", e => { e.UseRawJsonSerializer(); e.ConfigureConsumer<ProductConsumer>(context); e.ConfigureConsumer<NameConsumer>(context); }); cfg.DefaultContentType = new ContentType("application/json"); cfg.UseRawJsonDeserializer(); cfg.ConfigureEndpoints(context); }); });
NameConsumer.cs
public class NameConsumer: IConsumer<NameMessage> { private readonly ILogger<NameConsumer> _logger; public NameConsumer(ILogger<NameConsumer> logger) { _logger=logger; } public Task Consume(ConsumeContext<NameMessage> context) { _logger.LogInformation(" [NameConsumer] Message received context: {conetxt} ",JsonSerializer.Serialize(context)); return Task.CompletedTask; } } public interface NameMessage { public string Name { get; set; } }
ProductConsumer.cs
public class ProductConsumer: IConsumer<ProductMessage> { private readonly ILogger<ProductConsumer> _logger; public ProductConsumer(ILogger<ProductConsumer> logger) { _logger=logger; } public Task Consume(ConsumeContext<ProductMessage> context) { _logger.LogInformation(" [ProductConsumer] Message received context: {conetxt} ",JsonSerializer.Serialize(context)); return Task.CompletedTask; } } public interface ProductMessage { public string Code { get; set; } public string Name { get; set; } }
方案1:自定义消息头+消费者筛选
给消息添加自定义头(比如TargetConsumer),在消费者的消费逻辑里先检查头是否匹配,不匹配则直接跳过。
修改消费者代码
NameConsumer.cs
public Task Consume(ConsumeContext<NameMessage> context) { // 检查消息头是否匹配当前消费者 if (!context.Headers.TryGetHeader("TargetConsumer", out var target) || target.ToString() != nameof(NameConsumer)) { _logger.LogDebug("Skipping message not targeting NameConsumer"); return Task.CompletedTask; } _logger.LogInformation(" [NameConsumer] Message received context: {conetxt} ",JsonSerializer.Serialize(context)); return Task.CompletedTask; }
ProductConsumer.cs
public Task Consume(ConsumeContext<ProductMessage> context) { if (!context.Headers.TryGetHeader("TargetConsumer", out var target) || target.ToString() != nameof(ProductConsumer)) { _logger.LogDebug("Skipping message not targeting ProductConsumer"); return Task.CompletedTask; } _logger.LogInformation(" [ProductConsumer] Message received context: {conetxt} ",JsonSerializer.Serialize(context)); return Task.CompletedTask; }
发消息时:在RabbitMQ后台添加TargetConsumer头,值设为NameConsumer或ProductConsumer,对应消费者才会处理。
方案2:拆分独立接收队列
把两个消费者分配到不同的队列,这样从后台发消息时直接指定目标队列,自然只会触发对应消费者。
修改Program.cs
services.AddMassTransit(x => { x.AddConsumer<ProductConsumer>(); x.AddConsumer<NameConsumer>(); x.UsingRabbitMq((context, cfg) => { cfg.Host($"myRabbitHost.com", 5672, "VHOST" , configurator => { configurator.Username("_user"); configurator.Password("_pwd"); }); // 给ProductConsumer单独创建队列 cfg.ReceiveEndpoint("ProductQueue", e => { e.UseRawJsonSerializer(); e.ConfigureConsumer<ProductConsumer>(context); }); // 给NameConsumer单独创建队列 cfg.ReceiveEndpoint("NameQueue", e => { e.UseRawJsonSerializer(); e.ConfigureConsumer<NameConsumer>(context); }); cfg.DefaultContentType = new ContentType("application/json"); cfg.UseRawJsonDeserializer(); cfg.ConfigureEndpoints(context); }); });
发消息时:在RabbitMQ后台选择对应队列(ProductQueue或NameQueue)发布消息即可。
方案3:手动添加消息类型标识(模拟MassTransit类型匹配)
MassTransit默认通过信封里的MessageType字段识别消息类型,不用信封的话,可以手动给消息添加MessageType头,让MassTransit只匹配对应类型的消费者。
发消息时添加头
在RabbitMQ后台添加MessageType头,值为消息接口的完整类型名(比如YourNamespace.NameMessage或YourNamespace.ProductMessage)。
确保MassTransit识别类型
services.AddMassTransit(x => { x.AddConsumer<ProductConsumer>(); x.AddConsumer<NameConsumer>(); // 显式注册消息类型,帮助MassTransit识别 x.AddMessage<NameMessage>(); x.AddMessage<ProductMessage>(); x.UsingRabbitMq((context, cfg) => { cfg.Host($"myRabbitHost.com", 5672, "VHOST" , configurator => { configurator.Username("_user"); configurator.Password("_pwd"); }); cfg.ReceiveEndpoint("AUTO.QUEUENAME", e => { e.UseRawJsonSerializer(); e.ConfigureConsumer<ProductConsumer>(context); e.ConfigureConsumer<NameConsumer>(context); // 启用基于消息类型头的筛选 e.UseMessageTypeRouting(); }); cfg.DefaultContentType = new ContentType("application/json"); cfg.UseRawJsonDeserializer(); cfg.ConfigureEndpoints(context); }); });
这样当你发消息时指定MessageType头为NameMessage的完整类型名,只有NameConsumer会触发;指定ProductMessage则只有ProductConsumer触发。
内容的提问来源于stack exchange,提问作者EvaHHHH

