C#中如何为所有DAPR outgoing消息自动注入CorrelationId元数据?
为所有Dapr Outgoing消息自动注入CorrelationId元数据
问题场景
已实现CorrelationId中间件处理常规HTTP请求的CorrelationId传递,但DaprClient为单例注册,每次调用PublishEventAsync都需手动构造并传递包含CorrelationId的元数据字典,希望实现自动注入,消除重复代码。使用版本:Dapr.AspNetCore v1.11.0、Dapr.Workflow v1.11.0。
解决方案
方案一:装饰器模式包装DaprClient
通过实现IDaprClient装饰器,在PublishEventAsync方法中自动合并CorrelationId元数据,同时保留原有DaprClient的全部功能。
1. 实现装饰器类
public class CorrelationIdDaprClientDecorator : IDaprClient { private readonly IDaprClient _innerClient; private readonly IServiceScopeFactory _scopeFactory; public CorrelationIdDaprClientDecorator(IDaprClient innerClient, IServiceScopeFactory scopeFactory) { _innerClient = innerClient; _scopeFactory = scopeFactory; } // 重写PublishEventAsync核心方法,自动注入CorrelationId public Task PublishEventAsync(string pubsubName, string topicName, object data, Dictionary<string, string> metadata = null, CancellationToken cancellationToken = default) { var correlationId = GetCurrentCorrelationId(); var mergedMetadata = MergeMetadata(metadata, correlationId); return _innerClient.PublishEventAsync(pubsubName, topicName, data, mergedMetadata, cancellationToken); } public Task PublishEventAsync<TData>(string pubsubName, string topicName, TData data, Dictionary<string, string> metadata = null, CancellationToken cancellationToken = default) { var correlationId = GetCurrentCorrelationId(); var mergedMetadata = MergeMetadata(metadata, correlationId); return _innerClient.PublishEventAsync(pubsubName, topicName, data, mergedMetadata, cancellationToken); } // 获取当前作用域的CorrelationId private string GetCurrentCorrelationId() { using var scope = _scopeFactory.CreateScope(); var generator = scope.ServiceProvider.GetRequiredService<ICorrelationIdGenerator>(); return generator.Get() ?? Guid.NewGuid().ToString(); } // 合并现有元数据与CorrelationId private Dictionary<string, string> MergeMetadata(Dictionary<string, string> existingMetadata, string correlationId) { var metadata = existingMetadata ?? new Dictionary<string, string>(); if (!metadata.ContainsKey(Http.Headers.CORRELATION_ID)) { metadata.Add(Http.Headers.CORRELATION_ID, correlationId); } return metadata; } // 转发IDaprClient其他所有方法到内部实例 public Task<TResponse> InvokeMethodAsync<TResponse>(HttpMethod httpMethod, string appId, string methodName, object data = null, Dictionary<string, string> metadata = null, CancellationToken cancellationToken = default) => _innerClient.InvokeMethodAsync<TResponse>(httpMethod, appId, methodName, data, metadata, cancellationToken); // 省略其他IDaprClient接口方法的实现,全部直接转发即可 }
2. 修改服务注册
替换默认的DaprClient为装饰器实例,解决单例与Scoped服务的生命周期冲突:
Services.AddScoped<ICorrelationIdGenerator, CorrelationIdGenerator>(); // 注册原始DaprClient Services.AddDaprClient(); // 用装饰器包装DaprClient Services.AddSingleton<IDaprClient>(sp => { var innerDaprClient = sp.GetRequiredService<DaprClient>(); var scopeFactory = sp.GetRequiredService<IServiceScopeFactory>(); return new CorrelationIdDaprClientDecorator(innerDaprClient, scopeFactory); });
方案二:自定义HTTP消息处理程序
利用DaprClient底层的HttpClient管道,添加自定义Handler自动注入CorrelationId元数据(适用于需要给所有Dapr HTTP请求添加CorrelationId的场景)。
1. 实现DelegatingHandler
public class CorrelationIdDaprHandler : DelegatingHandler { private readonly IServiceScopeFactory _scopeFactory; public CorrelationIdDaprHandler(IServiceScopeFactory scopeFactory) { _scopeFactory = scopeFactory; } protected override async Task<HttpResponseMessage> SendAsync(HttpRequestMessage request, CancellationToken cancellationToken) { using var scope = _scopeFactory.CreateScope(); var generator = scope.ServiceProvider.GetRequiredService<ICorrelationIdGenerator>(); var correlationId = generator.Get() ?? Guid.NewGuid().ToString(); // DaprClient会将metadata转换为带dapr-metadata-前缀的请求头 var metadataHeaderKey = $"dapr-metadata-{Http.Headers.CORRELATION_ID}"; if (!request.Headers.Contains(metadataHeaderKey)) { request.Headers.Add(metadataHeaderKey, correlationId); } return await base.SendAsync(request, cancellationToken); } }
2. 修改服务注册
配置DaprClient使用自定义Handler:
Services.AddScoped<ICorrelationIdGenerator, CorrelationIdGenerator>(); Services.AddDaprClient(builder => { builder.UseHttpMessageHandler(sp => { var handler = new CorrelationIdDaprHandler(sp.GetRequiredService<IServiceScopeFactory>()); handler.InnerHandler = new HttpClientHandler(); return handler; }); });
使用说明
两种方案实现后,调用PublishEventAsync时无需手动传递metadata,CorrelationId会自动注入:
// 直接调用,无需手动构造metadata await _client.PublishEventAsync("component", "my-topic", data);
内容的提问来源于stack exchange,提问作者Serhat
相关产品推荐
相关产品推荐

