如何在C#依赖注入作用域中添加消息特定的MyUserInfo实例?
实现方案
完全可以实现,核心思路是利用.NET的**作用域(Scope)**机制,在每个消息处理流程中创建独立的作用域,并将提取到的MyUserInfo存入该作用域内的上下文容器中,让Scoped服务通过依赖注入获取。这样既复用了根容器的服务注册,又能保证用户信息仅在当前消息处理的作用域内有效,和Web应用的请求作用域逻辑一致。
具体步骤
1. 定义Scoped上下文容器
创建一个用来存储当前用户信息的Scoped服务,每个作用域内会有独立的实例:
public class MyUserContext { // 存储当前消息中的用户信息 public MyUserInfo CurrentUser { get; set; } }
2. 注册服务
在Program.cs中将上下文容器和你的业务服务注册为Scoped:
var builder = WebApplication.CreateBuilder(args); // 注册Scoped的用户上下文容器 builder.Services.AddScoped<MyUserContext>(); // 注册你的业务Scoped服务 builder.Services.AddScoped<MyService>(); builder.Services.AddScoped<MyRepository>(); // 注册消息处理服务(比如Worker服务) builder.Services.AddHostedService<MessageProcessor>(); var app = builder.Build(); app.Run();
3. 消息处理时注入用户信息
在消息处理的入口逻辑中,通过IServiceScopeFactory创建作用域,提取用户信息并赋值到上下文容器,再从该作用域获取业务服务处理消息:
public class MessageProcessor : BackgroundService { private readonly IServiceScopeFactory _scopeFactory; // 假设你使用Azure Service Bus的客户端,其他SDK逻辑类似 private readonly ServiceBusProcessor _processor; public MessageProcessor(IServiceScopeFactory scopeFactory, ServiceBusClient client) { _scopeFactory = scopeFactory; _processor = client.CreateProcessor("your-queue-name", new ServiceBusProcessorOptions()); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _processor.ProcessMessageAsync += ProcessMessageHandler; _processor.ProcessErrorAsync += ProcessErrorHandler; await _processor.StartProcessingAsync(stoppingToken); await Task.Delay(Timeout.Infinite, stoppingToken); } private async Task ProcessMessageHandler(ProcessMessageEventArgs args) { // 创建当前消息处理的独立作用域 using var scope = _scopeFactory.CreateScope(); var userContext = scope.ServiceProvider.GetRequiredService<MyUserContext>(); // 从消息中提取MyUserInfo(根据实际消息格式调整) userContext.CurrentUser = ExtractUserInfo(args.Message); // 从当前作用域获取业务服务处理消息 var myService = scope.ServiceProvider.GetRequiredService<MyService>(); await myService.ProcessMessage(args.Message, args.CancellationToken); // 完成消息处理 await args.CompleteMessageAsync(args.Message); } private MyUserInfo ExtractUserInfo(ServiceBusReceivedMessage message) { // 示例:从消息应用属性中提取用户ID if (message.ApplicationProperties.TryGetValue("UserId", out var userIdObj) && userIdObj is string userId) { return new MyUserInfo { UserId = userId, /* 其他属性 */ }; } throw new InvalidOperationException("消息中未包含用户信息"); } private Task ProcessErrorHandler(ProcessErrorEventArgs args) { // 处理消息异常逻辑 Console.WriteLine($"消息处理出错: {args.Exception.Message}"); return Task.CompletedTask; } public override async Task StopAsync(CancellationToken cancellationToken) { await _processor.StopProcessingAsync(cancellationToken); await base.StopAsync(cancellationToken); } }
4. 在业务服务中使用用户信息
让你的Scoped业务服务通过依赖注入获取MyUserContext,直接读取当前作用域内的用户信息:
public class MyService { private readonly MyUserContext _userContext; private readonly MyRepository _repository; public MyService(MyUserContext userContext, MyRepository repository) { _userContext = userContext; _repository = repository; } public async Task ProcessMessage(ServiceBusReceivedMessage message, CancellationToken cancellationToken) { // 获取当前作用域的用户信息 var currentUser = _userContext.CurrentUser; if (currentUser == null) { throw new InvalidOperationException("当前作用域无用户信息"); } // 执行业务逻辑,比如调用Repository await _repository.LogUserAction(currentUser.UserId, message.MessageId, cancellationToken); } }
方案优势
- 作用域隔离:每个消息处理对应独立的Scope,用户信息仅在当前Scope内有效,不会干扰其他消息处理。
- 复用根容器:无需新建
ServiceProvider,通过IServiceScopeFactory创建Scope,复用根容器的服务注册,性能更优。 - 和Web请求作用域一致:逻辑和Web应用中请求作用域传递用户信息的模式完全对齐,符合.NET的依赖注入设计规范。
注意事项
- 所有需要访问用户信息的业务逻辑,必须从当前消息处理的Scope中获取服务,不能直接使用Singleton服务(除非Singleton服务通过
IServiceScopeFactory获取Scope内的上下文,但不推荐)。 - 确保Scope被正确释放:使用
using包裹Scope,保证处理完成后自动释放所有Scoped服务。 - 如果使用其他服务总线SDK(如RabbitMQ、Kafka),只需调整消息接收和处理的代码,作用域管理的逻辑完全通用。
内容的提问来源于stack exchange,提问作者DLeh
相关产品推荐
相关产品推荐

