You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在.NET 5.0+CQRS+PostgreSQL主从架构(K8s容器化)中配置MediatR实现读写分离连接字符串

如何为基于MediatR的CQRS配置独立的PostgreSQL读写连接字符串?

我们的项目基于.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.30 15:07:41