MassTransit中Send Filter无法被IRequestClient触发的配置问题
问题场景
需要捕获HTTP请求携带的元数据(包括user-id、trace-id、tenant-id等),通过MassTransit的IRequestClient实现元数据跨服务透传。
发送消息的客户端基于HotChocolate实现GraphQL API,代码如下:
public async Task<Client> ClientCreate( [Argument] ClientCreateInput input, [Service] IMapper mapper, [Service] IRequestClient<ClientCreate> requester, [Service] IMessageValidator<ClientCreate> validator, CancellationToken cancellationToken) { var command = await validator.ValidateAsync( mapper.Map<ClientCreate>(input), cancellationToken); var response = await requester.GetResponse<ClientContract, ErrorResponse>( command, cancellationToken); return response.TransformResponseMessage(mapper.Map<Client>); }
为实现元数据透传,编写了如下发送过滤器:
public class SendApiDomainMetadataFilter : IFilter<SendContext> { public void Probe(ProbeContext context) { context.CreateFilterScope(nameof(SendApiDomainMetadataFilter)); } public async Task Send(SendContext context, IPipe<SendContext> next) { Console.WriteLine("Sending metadata ..."); if (context.TryGetPayload(out IServiceProvider? serviceProvider) && serviceProvider!.GetService<IResolverContext>() is { } resolverContext && resolverContext.GetDomainMetadata() is {} metadata) { var serializedMetadata = System.Text.Json.JsonSerializer.Serialize(metadata); Console.WriteLine("Sending metadata: '{0}'", serializedMetadata); context.Headers.Set(nameof(DomainMetadata), serializedMetadata); } await next.Send(context); } }
遇到的问题
该过滤器始终不会触发执行,尝试了两种配置方式均无效:
builder.Services.AddMassTransit(bus => { bus.UsingRabbitMq((context, configurator) => { // ... 其他基础配置 configurator.UseConsumeFilter(typeof(ValidateConsumeFilter<>), context); // 第一种配置方式 configurator.UseSendFilter(typeof(SendApiDomainMetadataFilter), context); // 第二种配置方式 configurator.ConfigureSend(sendConfigurator => { Console.WriteLine("Configuring send pipeline"); sendConfigurator.UseFilter(new SendApiDomainMetadataFilter()); }); configurator.ConfigureEndpoints(context); configurator.Host(...); }); });
最初以为configurator.UseSendFilter和sendConfigurator.UseFilter二选一即可让过滤器生效,但两种配置都没达到预期:所有IRequestClient.Send调用发出的消息,在发送到RabbitMQ前都没有经过SendApiDomainMetadataFilter.Send方法处理。
问题排查过程
参照已有的消费端过滤器实现,尝试编写作用域过滤器,但公开的示例全部针对消费端,没有发送端在非消费者/非Saga场景下的配置参考。排查中发现核心现象:从消费者、Saga之外的位置发送消息时,Send管道不会被调用。
编写了一个泛型测试过滤器验证问题:
public class DummySendFilter<TMessage> : IFilter<SendContext<TMessage>> where TMessage : class { public void Probe(ProbeContext context) { context.CreateFilterScope("DummySendFilter"); } public async Task Send(SendContext<TMessage> context, IPipe<SendContext<TMessage>> next) { Console.WriteLine("Dummy Sending " + typeof(TMessage).Name); await next.Send(context); } }
架构中网关服务、业务服务两个服务通过MassTransit通信,在两个服务中都添加了该测试过滤器,配置如下:
// 位于UsingRabbitMq配置块内部 configurator.UseSendFilter(typeof(DummySendFilter<>), svp);
测试结果:
- 网关通过
IRequestClient<ClientCreate>发送请求消息时,不会经过DummySendFilter,消息直接发送到RabbitMQ - 消息到达业务服务消费者后,会正常经过验证过滤器
- 业务服务处理完成返回响应消息时,会正常经过配置的DummySendFilter,打印
Dummy Sending ClientContract日志
网关侧的发送过滤器始终无法激活,无法在消息到达RabbitMQ前做自定义处理。
解决方案
问题根因是IRequestClient发起请求时,并非始终走Send管道,部分场景下会走Publish管道,仅实现IFilter<SendContext>/IFilter<SendContext<TMessage>>的过滤器无法捕获这类请求消息。需要实现针对PublishContext的过滤器,或编写通用过滤器同时覆盖Send、Publish两个管道场景,配置后即可在网关侧拦截所有IRequestClient发出的消息,完成元数据注入。
按照该思路调整过滤器实现后,问题已解决。
内容的提问来源于stack exchange,提问作者isierra

