如何在.NET 5.0+CQRS+PostgreSQL主从架构(K8s容器化)中配置MediatR实现读写分离连接字符串
我们的项目基于.NET 5.0框架,采用CQRS架构并使用PostgreSQL数据库,所有应用均已容器化并部署在Kubernetes集群中。PostgreSQL采用主从架构,写入操作通过pgpool路由至主节点,读取操作直接访问从节点(pgpool配置在主节点之前)。请问是否可以为基于MediatR实现的CQRS配置独立的读写连接字符串?当前我们的依赖注入(DI)配置如下:
services.AddDbContext<DataContext>(opt => { opt.EnableDetailedErrors(); opt.UseLazyLoadingProxies(); opt.UseNpgsql(configuration.GetConnectionString("TaikunConnection"), options => { options.MigrationsAssembly(typeof(DataContext).Assembly.FullName); options.EnableRetryOnFailure(); options.CommandTimeout((int)TimeSpan.FromMinutes(10).TotalSeconds); }); }); services.AddScoped<IDataContext>(provider => provider.GetService<DataContext>()); var builder = services.AddIdentityCore<User>() .AddEntityFrameworkStores<DataContext>(); var identityBuilder = new IdentityBuilder(builder.UserType, builder.Services); identityBuilder.AddSignInManager<SignInManager<User>>(); identityBuilder.AddUserManager<UserManager<User>>();
当然可以!结合你当前的.NET 5 + MediatR + PostgreSQL主从架构,完全可以通过区分CQRS的命令(写)和查询(读)操作来配置独立的读写连接字符串,这样能完美适配pgpool的路由策略,提升系统性能。下面是一步步的实现方案:
1. 配置读写分离的连接字符串
首先在appsettings.json中分别定义主库(写操作)和从库(读操作)的连接字符串,对应你的pgpool主节点和从节点地址:
"ConnectionStrings": { "TaikunWriteConnection": "Host=pgpool-main;Database=taikun;Username=xxx;Password=xxx", "TaikunReadConnection": "Host=pgpool-replica;Database=taikun;Username=xxx;Password=xxx" }
2. 定义读写专属的DbContext
为了清晰划分CQRS边界,推荐创建两个DbContext:WriteDataContext(处理写操作)和ReadDataContext(处理读操作),它们可以继承自你现有的DataContext:
public class WriteDataContext : DataContext { public WriteDataContext(DbContextOptions<WriteDataContext> options) : base(options) { } } public class ReadDataContext : DataContext { public ReadDataContext(DbContextOptions<ReadDataContext> options) : base(options) { } }
这样做的好处是让代码意图更明确,避免混用读写连接。
3. 更新依赖注入配置
修改你的DI注册代码,分别为两个DbContext配置对应的连接字符串,同时注意Identity相关操作都是写操作,要绑定到WriteDataContext:
// 注册写操作DbContext(关联主库连接) services.AddDbContext<WriteDataContext>(opt => { opt.EnableDetailedErrors(); opt.UseLazyLoadingProxies(); opt.UseNpgsql(configuration.GetConnectionString("TaikunWriteConnection"), options => { options.MigrationsAssembly(typeof(DataContext).Assembly.FullName); options.EnableRetryOnFailure(); options.CommandTimeout((int)TimeSpan.FromMinutes(10).TotalSeconds); }); }); // 注册读操作DbContext(关联从库连接) services.AddDbContext<ReadDataContext>(opt => { opt.EnableDetailedErrors(); opt.UseLazyLoadingProxies(); opt.UseNpgsql(configuration.GetConnectionString("TaikunReadConnection"), options => { options.MigrationsAssembly(typeof(DataContext).Assembly.FullName); options.EnableRetryOnFailure(); options.CommandTimeout((int)TimeSpan.FromMinutes(10).TotalSeconds); }); }); // Identity服务绑定写DbContext,因为用户创建/更新都是写操作 var builder = services.AddIdentityCore<User>() .AddEntityFrameworkStores<WriteDataContext>(); var identityBuilder = new IdentityBuilder(builder.UserType, builder.Services); identityBuilder.AddSignInManager<SignInManager<User>>(); identityBuilder.AddUserManager<UserManager<User>>();
4. 在MediatR处理程序中注入对应DbContext
现在只需要在命令处理程序(写操作)中注入WriteDataContext,查询处理程序(读操作)中注入ReadDataContext即可:
命令处理程序示例(写操作)
public class CreateUserCommandHandler : IRequestHandler<CreateUserCommand, bool> { private readonly WriteDataContext _writeDbContext; public CreateUserCommandHandler(WriteDataContext writeDbContext) { _writeDbContext = writeDbContext; } public async Task<bool> Handle(CreateUserCommand request, CancellationToken cancellationToken) { // 所有写操作都走主库连接 _writeDbContext.Users.Add(new User { UserName = request.UserName, Email = request.Email }); await _writeDbContext.SaveChangesAsync(cancellationToken); return true; } }
查询处理程序示例(读操作)
public class GetUserQueryHandler : IRequestHandler<GetUserQuery, UserDto> { private readonly ReadDataContext _readDbContext; public GetUserQueryHandler(ReadDataContext readDbContext) { _readDbContext = readDbContext; } public async Task<UserDto> Handle(GetUserQuery request, CancellationToken cancellationToken) { // 所有读操作都走从库连接 var user = await _readDbContext.Users .FirstOrDefaultAsync(u => u.Id == request.UserId, cancellationToken); return user == null ? null : new UserDto { Id = user.Id, UserName = user.UserName, Email = user.Email }; } }
5. 可选优化:自动匹配读写DbContext的MediatR行为
如果不想每个处理程序都手动区分注入的DbContext,可以创建一个MediatR管道行为,根据请求类型(Command/Query)自动切换对应的DbContext:
public class CqrsDbContextBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse> where TRequest : IRequest<TResponse> { private readonly WriteDataContext _writeDbContext; private readonly ReadDataContext _readDbContext; private readonly IServiceScopeFactory _scopeFactory; public CqrsDbContextBehavior(WriteDataContext writeDbContext, ReadDataContext readDbContext, IServiceScopeFactory scopeFactory) { _writeDbContext = writeDbContext; _readDbContext = readDbContext; _scopeFactory = scopeFactory; } public async Task<TResponse> Handle(TRequest request, RequestHandlerDelegate<TResponse> next, CancellationToken cancellationToken) { // 根据请求名称后缀判断是命令还是查询 var isCommand = typeof(TRequest).Name.EndsWith("Command"); var targetDbContext = isCommand ? _writeDbContext : _readDbContext; // 在当前作用域中替换IDataContext为目标实例(如果你的处理程序依赖IDataContext) using var scope = _scopeFactory.CreateScope(); var scopeServices = scope.ServiceProvider; // 这里可以通过自定义DI扩展逻辑绑定targetDbContext到IDataContext,示例仅作参考 // 更简单的方式是让处理程序直接依赖Write/ReadDataContext,避免替换的麻烦 return await next(); } }
然后注册这个行为到DI:
services.AddTransient(typeof(IPipelineBehavior<,>), typeof(CqrsDbContextBehavior<,>));
这样就能实现完全的读写分离,让你的CQRS架构完美适配PostgreSQL主从集群的路由策略了。
内容的提问来源于stack exchange,提问作者Arzu Suleymanov

