📖 数据密集型设计

数据库读写分离与主从复制

深入探讨数据库读写分离架构、主从复制原理及一致性保障

一、读写分离概述

读写分离是数据库水平扩展的常用策略,通过将读操作和写操作分离到不同的数据库节点,能够提升系统的吞吐量和可用性。主从复制是实现读写分离的基础,通过异步或半同步复制,将主库的数据同步到从库。

二、主从复制原理

2.1 MySQL主从复制架构

graph TD A[客户端] --> B[应用层] B --> C{SQL类型} C -->|写操作| D[主库Master] C -->|读操作| E[负载均衡] D --> F[Binlog记录] F --> G[从库Slave1] F --> H[从库Slave2] F --> I[从库Slave3] E --> G E --> H E --> I G --> G1[Relay Log] G1 --> G2[SQL线程执行] H --> H1[Relay Log] H1 --> H2[SQL线程执行] I --> I1[Relay Log] I1 --> I2[SQL线程执行]

2.2 复制类型

类型 描述 优点 缺点 适用场景
异步复制 主库写入后立即返回 性能好 可能丢失数据 非关键业务
半同步复制 等待至少一个从库确认 数据安全 性能稍有下降 关键业务
全同步复制 等待所有从库确认 数据最安全 性能差 金融级业务

三、读写分离实现

3.1 应用层读写分离

public class ReadWriteSplitDbContext : DbContext
{
    private readonly IDbConnectionFactory _masterConnectionFactory;
    private readonly IDbConnectionFactory _slaveConnectionFactory;
    private readonly IReadWriteSplitter _splitter;
    
    public ReadWriteSplitDbContext(
        IDbConnectionFactory masterConnectionFactory,
        IDbConnectionFactory slaveConnectionFactory,
        IReadWriteSplitter splitter)
    {
        _masterConnectionFactory = masterConnectionFactory;
        _slaveConnectionFactory = slaveConnectionFactory;
        _splitter = splitter;
    }
    
    protected override void OnConfiguring(DbContextOptionsBuilder optionsBuilder)
    {
        var connection = _splitter.UseMaster ? 
            _masterConnectionFactory.CreateConnection() : 
            _slaveConnectionFactory.CreateConnection();
        
        optionsBuilder.UseMySql(connection, ServerVersion.AutoDetect(connection));
    }
    
    public override int SaveChanges()
    {
        using (_splitter.SwitchToMaster())
        {
            return base.SaveChanges();
        }
    }
    
    public override async Task SaveChangesAsync(CancellationToken cancellationToken = default)
    {
        using (_splitter.SwitchToMaster())
        {
            return await base.SaveChangesAsync(cancellationToken);
        }
    }
}

public class ReadWriteSplitter : IReadWriteSplitter
{
    private static readonly AsyncLocal _useMaster = new();
    
    public bool UseMaster
    {
        get => _useMaster.Value;
        set => _useMaster.Value = value;
    }
    
    public IDisposable SwitchToMaster()
    {
        var previousValue = _useMaster.Value;
        _useMaster.Value = true;
        
        return new DisposableAction(() => _useMaster.Value = previousValue);
    }
    
    public IDisposable SwitchToSlave()
    {
        var previousValue = _useMaster.Value;
        _useMaster.Value = false;
        
        return new DisposableAction(() => _useMaster.Value = previousValue);
    }
}

public class DisposableAction : IDisposable
{
    private readonly Action _action;
    
    public DisposableAction(Action action) => _action = action;
    
    public void Dispose() => _action();
}

3.2 中间件读写分离

public class DbProxyMiddleware
{
    private readonly RequestDelegate _next;
    private readonly IDbRouter _dbRouter;
    
    public DbProxyMiddleware(RequestDelegate next, IDbRouter dbRouter)
    {
        _next = next;
        _dbRouter = dbRouter;
    }
    
    public async Task InvokeAsync(HttpContext context)
    {
        var route = _dbRouter.Route(context.Request);
        
        context.Items["DbRoute"] = route;
        
        await _next(context);
    }
}

public class DbRouter : IDbRouter
{
    private readonly List _slaveNodes;
    private readonly ILoadBalancer _loadBalancer;
    
    public DbRouter(List slaveNodes, ILoadBalancer loadBalancer)
    {
        _slaveNodes = slaveNodes;
        _loadBalancer = loadBalancer;
    }
    
    public DbRoute Route(HttpRequest request)
    {
        if (IsWriteRequest(request))
        {
            return new DbRoute(DbNodeType.Master, "master");
        }
        
        var slaveNode = _loadBalancer.Select(_slaveNodes);
        
        return new DbRoute(DbNodeType.Slave, slaveNode.ConnectionString);
    }
    
    private bool IsWriteRequest(HttpRequest request)
    {
        var writeMethods = new HashSet { "POST", "PUT", "DELETE", "PATCH" };
        
        return writeMethods.Contains(request.Method.ToUpper());
    }
    
    public DbRoute RouteForTransaction()
    {
        return new DbRoute(DbNodeType.Master, "master");
    }
    
    public DbRoute RouteForReadAfterWrite()
    {
        return new DbRoute(DbNodeType.Master, "master");
    }
}

public class DbRoute
{
    public DbNodeType NodeType { get; }
    public string ConnectionString { get; }
    
    public DbRoute(DbNodeType nodeType, string connectionString)
    {
        NodeType = nodeType;
        ConnectionString = connectionString;
    }
}

public enum DbNodeType { Master, Slave }

四、主从复制配置

4.1 MySQL主从配置

public class MySqlReplicationConfigService
{
    public async Task ConfigureMasterAsync(string masterHost, int masterPort, string masterUser, string masterPassword)
    {
        var config = new MySqlConfig
        {
            ServerId = 1,
            LogBin = "mysql-bin",
            BinlogFormat = BinlogFormat.Row,
            BinlogRowImage = BinlogRowImage.Full,
            ExpireLogsDays = 7
        };
        
        await ApplyConfigAsync(masterHost, masterPort, masterUser, masterPassword, config);
        
        await CreateReplicationUserAsync(masterHost, masterPort, masterUser, masterPassword, "repl", "repl_password");
    }
    
    public async Task ConfigureSlaveAsync(string slaveHost, int slavePort, string slaveUser, string slavePassword,
        string masterHost, int masterPort, string replUser, string replPassword)
    {
        var config = new MySqlConfig
        {
            ServerId = await GetNextServerIdAsync(),
            LogSlaveUpdates = true,
            ReadOnly = true,
            RelayLog = "relay-bin"
        };
        
        await ApplyConfigAsync(slaveHost, slavePort, slaveUser, slavePassword, config);
        
        await InitializeSlaveAsync(slaveHost, slavePort, slaveUser, slavePassword,
            masterHost, masterPort, replUser, replPassword);
    }
    
    private async Task ApplyConfigAsync(string host, int port, string user, string password, MySqlConfig config)
    {
        using var connection = new MySqlConnection($"server={host};port={port};user={user};password={password}");
        await connection.OpenAsync();
        
        await ExecuteConfigCommandAsync(connection, $"SET GLOBAL server_id = {config.ServerId}");
        await ExecuteConfigCommandAsync(connection, $"SET GLOBAL log_bin = '{config.LogBin}'");
        await ExecuteConfigCommandAsync(connection, $"SET GLOBAL binlog_format = '{config.BinlogFormat}'");
    }
    
    private async Task CreateReplicationUserAsync(string host, int port, string user, string password, 
        string replUser, string replPassword)
    {
        using var connection = new MySqlConnection($"server={host};port={port};user={user};password={password}");
        await connection.OpenAsync();
        
        var sql = $"CREATE USER '{replUser}'@'%' IDENTIFIED BY '{replPassword}'";
        await connection.ExecuteAsync(sql);
        
        sql = $"GRANT REPLICATION SLAVE ON *.* TO '{replUser}'@'%'";
        await connection.ExecuteAsync(sql);
    }
    
    private async Task InitializeSlaveAsync(string slaveHost, int slavePort, string slaveUser, string slavePassword,
        string masterHost, int masterPort, string replUser, string replPassword)
    {
        using var connection = new MySqlConnection($"server={slaveHost};port={slavePort};user={slaveUser};password={slavePassword}");
        await connection.OpenAsync();
        
        var sql = @$"CHANGE MASTER TO
            MASTER_HOST='{masterHost}',
            MASTER_PORT={masterPort},
            MASTER_USER='{replUser}',
            MASTER_PASSWORD='{replPassword}',
            MASTER_AUTO_POSITION=1";
        
        await connection.ExecuteAsync(sql);
        
        await connection.ExecuteAsync("START SLAVE");
    }
}

public class MySqlConfig
{
    public int ServerId { get; set; }
    public string LogBin { get; set; }
    public BinlogFormat BinlogFormat { get; set; }
    public BinlogRowImage BinlogRowImage { get; set; }
    public int ExpireLogsDays { get; set; }
    public bool LogSlaveUpdates { get; set; }
    public bool ReadOnly { get; set; }
    public string RelayLog { get; set; }
}

public enum BinlogFormat { Statement, Row, Mixed }
public enum BinlogRowImage { Full, Minimal, Noblob }

4.2 复制状态监控

public class ReplicationMonitorService
{
    public async Task GetReplicationStatusAsync(string slaveHost, int slavePort, string user, string password)
    {
        using var connection = new MySqlConnection($"server={slaveHost};port={slavePort};user={user};password={password}");
        await connection.OpenAsync();
        
        var reader = await connection.ExecuteReaderAsync("SHOW SLAVE STATUS");
        
        if (!reader.HasRows)
        {
            return new ReplicationStatus { IsRunning = false };
        }
        
        await reader.ReadAsync();
        
        return new ReplicationStatus
        {
            IsRunning = reader.GetString("Slave_IO_Running") == "Yes" && 
                        reader.GetString("Slave_SQL_Running") == "Yes",
            MasterHost = reader.GetString("Master_Host"),
            MasterPort = reader.GetInt32("Master_Port"),
            RelayLogFile = reader.GetString("Relay_Master_Log_File"),
            RelayLogPosition = reader.GetInt64("Exec_Master_Log_Pos"),
            MasterLogFile = reader.GetString("Master_Log_File"),
            MasterLogPosition = reader.GetInt64("Read_Master_Log_Pos"),
            SecondsBehindMaster = reader.GetInt32("Seconds_Behind_Master"),
            LastError = reader.GetString("Last_Error")
        };
    }
    
    public async Task> GetAllSlavesStatusAsync(List slaves)
    {
        var statuses = new List();
        
        foreach (var slave in slaves)
        {
            var status = await GetReplicationStatusAsync(slave.Host, slave.Port, slave.User, slave.Password);
            status.SlaveHost = slave.Host;
            statuses.Add(status);
        }
        
        return statuses;
    }
    
    public async Task CheckReplicationHealthAsync(List slaves, int maxLagSeconds = 30)
    {
        var statuses = await GetAllSlavesStatusAsync(slaves);
        
        return statuses.All(s => s.IsRunning && s.SecondsBehindMaster <= maxLagSeconds);
    }
    
    public async Task CheckAndAlertAsync(List slaves)
    {
        var statuses = await GetAllSlavesStatusAsync(slaves);
        
        var stoppedSlaves = statuses.Where(s => !s.IsRunning).ToList();
        var laggingSlaves = statuses.Where(s => s.IsRunning && s.SecondsBehindMaster > 30).ToList();
        
        if (stoppedSlaves.Count > 0 || laggingSlaves.Count > 0)
        {
            return new ReplicationAlert
            {
                AlertType = stoppedSlaves.Count > 0 ? AlertType.ReplicationStopped : AlertType.HighLag,
                StoppedSlaves = stoppedSlaves.Select(s => s.SlaveHost).ToList(),
                LaggingSlaves = laggingSlaves.Select(s => new { s.SlaveHost, s.SecondsBehindMaster }).ToList(),
                Timestamp = DateTime.UtcNow
            };
        }
        
        return null;
    }
}

public class ReplicationStatus
{
    public string SlaveHost { get; set; }
    public bool IsRunning { get; set; }
    public string MasterHost { get; set; }
    public int MasterPort { get; set; }
    public string RelayLogFile { get; set; }
    public long RelayLogPosition { get; set; }
    public string MasterLogFile { get; set; }
    public long MasterLogPosition { get; set; }
    public int SecondsBehindMaster { get; set; }
    public string LastError { get; set; }
}

public enum AlertType { ReplicationStopped, HighLag, ConnectionError }

五、一致性保障

5.1 读写一致性策略

graph TD A[写操作] --> B[主库写入] B --> C{一致性级别} C -->|强一致性| D[等待从库同步] D --> E[从库返回确认] E --> F[客户端返回成功] C -->|最终一致性| G[立即返回成功] G --> H[异步同步到从库] I[读操作] --> J{读策略} J -->|读自己写| K[强制读主库] J -->|读最新数据| L[检查从库延迟] L -->|延迟>阈值| K L -->|延迟<=阈值| M[读从库] J -->|读任意| M

5.2 读自己写实现

public class ReadYourWritesService
{
    private readonly IReadWriteSplitter _splitter;
    private readonly IDistributedCache _cache;
    
    public ReadYourWritesService(IReadWriteSplitter splitter, IDistributedCache cache)
    {
        _splitter = splitter;
        _cache = cache;
    }
    
    public async Task RecordWriteAsync(string userId, string entityType, string entityId)
    {
        var key = $"read-your-writes:{userId}:{entityType}:{entityId}";
        
        await _cache.SetStringAsync(key, "true", new DistributedCacheEntryOptions
        {
            AbsoluteExpirationRelativeToNow = TimeSpan.FromMinutes(5)
        });
    }
    
    public async Task ShouldReadMasterAsync(string userId, string entityType, string entityId)
    {
        var key = $"read-your-writes:{userId}:{entityType}:{entityId}";
        
        var value = await _cache.GetStringAsync(key);
        
        return !string.IsNullOrEmpty(value);
    }
    
    public async Task ExecuteWithReadYourWritesAsync(
        string userId, string entityType, string entityId,
        Func> readOperation)
    {
        var shouldReadMaster = await ShouldReadMasterAsync(userId, entityType, entityId);
        
        using (shouldReadMaster ? _splitter.SwitchToMaster() : _splitter.SwitchToSlave())
        {
            return await readOperation();
        }
    }
    
    public async Task ExecuteWriteAndReadAsync(
        string userId, string entityType, string entityId,
        Func writeOperation,
        Func> readOperation)
    {
        await writeOperation();
        
        await RecordWriteAsync(userId, entityType, entityId);
        
        return await ExecuteWithReadYourWritesAsync(userId, entityType, entityId, readOperation);
    }
}

5.3 延迟感知路由

public class LatencyAwareRouter : IDbRouter
{
    private readonly List _slaveNodes;
    private readonly IReplicationMonitorService _monitorService;
    private readonly int _maxLagThresholdSeconds;
    
    public LatencyAwareRouter(List slaveNodes, 
        IReplicationMonitorService monitorService,
        int maxLagThresholdSeconds = 30)
    {
        _slaveNodes = slaveNodes;
        _monitorService = monitorService;
        _maxLagThresholdSeconds = maxLagThresholdSeconds;
    }
    
    public async Task RouteAsync(HttpRequest request)
    {
        if (IsWriteRequest(request))
        {
            return new DbRoute(DbNodeType.Master, "master");
        }
        
        var statuses = await _monitorService.GetAllSlavesStatusAsync(_slaveNodes);
        
        var healthySlaves = statuses
            .Where(s => s.IsRunning && s.SecondsBehindMaster <= _maxLagThresholdSeconds)
            .ToList();
        
        if (healthySlaves.Count == 0)
        {
            return new DbRoute(DbNodeType.Master, "master");
        }
        
        var selectedSlave = SelectRandomSlave(healthySlaves);
        
        return new DbRoute(DbNodeType.Slave, selectedSlave.ConnectionString);
    }
    
    private SlaveNode SelectRandomSlave(List healthySlaves)
    {
        var random = new Random();
        var index = random.Next(healthySlaves.Count);
        var slaveHost = healthySlaves[index].SlaveHost;
        
        return _slaveNodes.First(s => s.Host == slaveHost);
    }
    
    private bool IsWriteRequest(HttpRequest request)
    {
        var writeMethods = new HashSet { "POST", "PUT", "DELETE", "PATCH" };
        
        return writeMethods.Contains(request.Method.ToUpper());
    }
}

六、故障转移

6.1 主库故障转移

public class MasterFailoverService
{
    private readonly List _slaveNodes;
    private readonly IReplicationMonitorService _monitorService;
    
    public MasterFailoverService(List slaveNodes, 
        IReplicationMonitorService monitorService)
    {
        _slaveNodes = slaveNodes;
        _monitorService = monitorService;
    }
    
    public async Task PerformFailoverAsync(string currentMasterHost)
    {
        var statuses = await _monitorService.GetAllSlavesStatusAsync(_slaveNodes);
        
        var candidate = statuses
            .Where(s => s.IsRunning && s.SlaveHost != currentMasterHost)
            .OrderBy(s => s.SecondsBehindMaster)
            .FirstOrDefault();
        
        if (candidate == null)
        {
            return new FailoverResult { Success = false, Message = "没有可用的从库候选" };
        }
        
        var slaveNode = _slaveNodes.First(s => s.Host == candidate.SlaveHost);
        
        await PromoteToMasterAsync(slaveNode);
        
        await ReconfigureOtherSlavesAsync(slaveNode, statuses);
        
        return new FailoverResult
        {
            Success = true,
            NewMasterHost = slaveNode.Host,
            Message = $"成功将 {slaveNode.Host} 提升为主库"
        };
    }
    
    private async Task PromoteToMasterAsync(SlaveNode slaveNode)
    {
        using var connection = new MySqlConnection(slaveNode.ConnectionString);
        await connection.OpenAsync();
        
        await connection.ExecuteAsync("STOP SLAVE");
        await connection.ExecuteAsync("RESET MASTER");
        await connection.ExecuteAsync("SET GLOBAL read_only = OFF");
    }
    
    private async Task ReconfigureOtherSlavesAsync(SlaveNode newMaster, List statuses)
    {
        foreach (var status in statuses)
        {
            if (status.SlaveHost == newMaster.Host)
            {
                continue;
            }
            
            var slaveNode = _slaveNodes.First(s => s.Host == status.SlaveHost);
            
            using var connection = new MySqlConnection(slaveNode.ConnectionString);
            await connection.OpenAsync();
            
            var sql = @$"CHANGE MASTER TO
                MASTER_HOST='{newMaster.Host}',
                MASTER_PORT={newMaster.Port},
                MASTER_USER='{newMaster.ReplUser}',
                MASTER_PASSWORD='{newMaster.ReplPassword}',
                MASTER_AUTO_POSITION=1";
            
            await connection.ExecuteAsync(sql);
            await connection.ExecuteAsync("START SLAVE");
        }
    }
}

public class FailoverResult
{
    public bool Success { get; set; }
    public string NewMasterHost { get; set; }
    public string Message { get; set; }
}

七、读写分离最佳实践

7.1 读写分离策略

策略 描述 实现方式 适用场景
按方法类型 POST/PUT/DELETE写,GET读 中间件拦截 RESTful API
按事务 事务内强制主库 事务拦截器 事务操作
读自己写 刚写入的数据读主库 缓存标记 实时性要求
延迟感知 根据从库延迟路由 监控+路由 高一致性要求

7.2 主从复制最佳实践

  • 使用Row格式的Binlog
  • 开启半同步复制
  • 监控复制延迟
  • 定期备份从库
  • 准备故障转移方案

八、总结

数据库读写分离与主从复制是提升数据库吞吐量的关键技术。通过将读操作分散到多个从库,能够显著提升系统的读性能。主从复制是实现读写分离的基础,需要根据业务需求选择合适的复制类型。一致性保障和故障转移是生产环境中必须考虑的问题,读自己写和延迟感知路由能够有效解决一致性问题。