📖 数据密集型设计

数据库分片与水平扩展

深入探讨数据库分片策略、水平扩展方案及分布式事务处理

一、数据库分片概述

数据库分片是将单一数据库拆分为多个数据库的过程,通过水平扩展能够处理海量数据和高并发访问。分片策略的选择直接影响系统的性能、扩展性和可维护性。

二、分片策略

2.1 分片策略对比

策略 描述 优点 缺点 适用场景
范围分片 按范围划分数据 易于扩展、范围查询高效 数据分布不均 时间序列、ID范围
哈希分片 按哈希值划分 数据均匀、查询简单 范围查询困难 用户ID、订单ID
列表分片 按列表值划分 灵活、业务友好 维护成本高 地区、分类
复合分片 多种策略组合 灵活、兼顾多种场景 复杂度高 复杂业务场景

2.2 分片架构

graph TD A[客户端] --> B[分片路由层] B --> C{分片键} C -->|范围分片| D[范围路由] C -->|哈希分片| E[哈希路由] C -->|列表分片| F[列表路由] D --> G[分片0] D --> H[分片1] D --> I[分片2] E --> G E --> H E --> I F --> G F --> H F --> I G --> G1[主库] G1 --> G2[从库1] G1 --> G3[从库2] H --> H1[主库] H1 --> H2[从库1] H1 --> H3[从库2] I --> I1[主库] I1 --> I2[从库1] I1 --> I3[从库2]

三、分片实现

3.1 分片路由实现

public class ShardingRouter
{
    private readonly Dictionary _strategies = new();
    
    public void RegisterStrategy(string tableName, IShardingStrategy strategy)
    {
        _strategies[tableName] = strategy;
    }
    
    public ShardingRoute Route(string tableName, object shardKey)
    {
        if (!_strategies.TryGetValue(tableName, out var strategy))
        {
            throw new InvalidOperationException($"未找到表 {tableName} 的分片策略");
        }
        
        return strategy.Route(shardKey);
    }
    
    public List RouteMultiple(string tableName, List shardKeys)
    {
        var routes = new HashSet();
        
        foreach (var key in shardKeys)
        {
            routes.Add(Route(tableName, key));
        }
        
        return routes.ToList();
    }
    
    public List RouteRange(string tableName, object startKey, object endKey)
    {
        if (!_strategies.TryGetValue(tableName, out var strategy))
        {
            throw new InvalidOperationException($"未找到表 {tableName} 的分片策略");
        }
        
        if (strategy is IRangeShardingStrategy rangeStrategy)
        {
            return rangeStrategy.RouteRange(startKey, endKey);
        }
        
        throw new InvalidOperationException($"表 {tableName} 不支持范围查询");
    }
}

public interface IShardingStrategy
{
    ShardingRoute Route(object shardKey);
}

public interface IRangeShardingStrategy : IShardingStrategy
{
    List RouteRange(object startKey, object endKey);
}

public class ShardingRoute
{
    public string DatabaseName { get; set; }
    public string TableName { get; set; }
}
            
            

3.2 哈希分片策略

public class HashShardingStrategy : IShardingStrategy
{
    private readonly int _databaseCount;
    private readonly int _tableCount;
    private readonly string _databasePrefix;
    private readonly string _tablePrefix;
    
    public HashShardingStrategy(int databaseCount, int tableCount, 
        string databasePrefix = "db_", string tablePrefix = "t_")
    {
        _databaseCount = databaseCount;
        _tableCount = tableCount;
        _databasePrefix = databasePrefix;
        _tablePrefix = tablePrefix;
    }
    
    public ShardingRoute Route(object shardKey)
    {
        var hash = ComputeHash(shardKey);
        
        var databaseIndex = hash % _databaseCount;
        var tableIndex = hash / _databaseCount % _tableCount;
        
        return new ShardingRoute
        {
            DatabaseName = $"{_databasePrefix}{databaseIndex}",
            TableName = $"{_tablePrefix}{tableIndex}"
        };
    }
    
    private int ComputeHash(object shardKey)
    {
        if (shardKey is string strKey)
        {
            return Math.Abs(strKey.GetHashCode());
        }
        
        if (shardKey is int intKey)
        {
            return intKey;
        }
        
        if (shardKey is long longKey)
        {
            return (int)(longKey & int.MaxValue);
        }
        
        return Math.Abs(shardKey.GetHashCode());
    }
}

public class ConsistentHashShardingStrategy : IShardingStrategy
{
    private readonly ConsistentHashRing _hashRing;
    
    public ConsistentHashShardingStrategy(List nodes, int virtualNodeCount = 100)
    {
        _hashRing = new ConsistentHashRing(virtualNodeCount);
        
        foreach (var node in nodes)
        {
            _hashRing.AddNode(node);
        }
    }
    
    public ShardingRoute Route(object shardKey)
    {
        var node = _hashRing.GetNode(shardKey.ToString());
        
        return new ShardingRoute
        {
            DatabaseName = node,
            TableName = "default"
        };
    }
}

public class ConsistentHashRing
{
    private readonly SortedDictionary _ring = new();
    private readonly int _virtualNodeCount;
    
    public ConsistentHashRing(int virtualNodeCount)
    {
        _virtualNodeCount = virtualNodeCount;
    }
    
    public void AddNode(string node)
    {
        for (int i = 0; i < _virtualNodeCount; i++)
        {
            var hash = ComputeHash($"{node}-{i}");
            _ring[hash] = node;
        }
    }
    
    public void RemoveNode(string node)
    {
        for (int i = 0; i < _virtualNodeCount; i++)
        {
            var hash = ComputeHash($"{node}-{i}");
            _ring.Remove(hash);
        }
    }
    
    public string GetNode(string key)
    {
        if (_ring.Count == 0)
        {
            throw new InvalidOperationException("哈希环为空");
        }
        
        var hash = ComputeHash(key);
        
        foreach (var pair in _ring)
        {
            if (pair.Key >= hash)
            {
                return pair.Value;
            }
        }
        
        return _ring.First().Value;
    }
    
    private int ComputeHash(string key)
    {
        return Math.Abs(key.GetHashCode());
    }
}

3.3 范围分片策略

public class RangeShardingStrategy : IRangeShardingStrategy
{
    private readonly List _partitions;
    private readonly string _databasePrefix;
    private readonly string _tablePrefix;
    
    public RangeShardingStrategy(List partitions, 
        string databasePrefix = "db_", string tablePrefix = "t_")
    {
        _partitions = partitions.OrderBy(p => p.StartValue).ToList();
        _databasePrefix = databasePrefix;
        _tablePrefix = tablePrefix;
    }
    
    public ShardingRoute Route(object shardKey)
    {
        var partition = FindPartition(shardKey);
        
        return new ShardingRoute
        {
            DatabaseName = $"{_databasePrefix}{partition.DatabaseIndex}",
            TableName = $"{_tablePrefix}{partition.TableIndex}"
        };
    }
    
    public List RouteRange(object startKey, object endKey)
    {
        var routes = new HashSet();
        
        foreach (var partition in _partitions)
        {
            if (IsOverlapping(partition, startKey, endKey))
            {
                routes.Add(new ShardingRoute
                {
                    DatabaseName = $"{_databasePrefix}{partition.DatabaseIndex}",
                    TableName = $"{_tablePrefix}{partition.TableIndex}"
                });
            }
        }
        
        return routes.ToList();
    }
    
    private RangePartition FindPartition(object shardKey)
    {
        foreach (var partition in _partitions)
        {
            if (Compare(shardKey, partition.StartValue) >= 0 && 
                (partition.EndValue == null || Compare(shardKey, partition.EndValue) < 0))
            {
                return partition;
            }
        }
        
        throw new InvalidOperationException("未找到对应的分片");
    }
    
    private bool IsOverlapping(RangePartition partition, object startKey, object endKey)
    {
        return Compare(startKey, partition.EndValue ?? endKey) < 0 && 
               Compare(endKey, partition.StartValue) > 0;
    }
    
    private int Compare(object a, object b)
    {
        return Comparer.Default.Compare(a, b);
    }
}

public class RangePartition
{
    public object StartValue { get; set; }
    public object EndValue { get; set; }
    public int DatabaseIndex { get; set; }
    public int TableIndex { get; set; }
}
            
            

四、分片管理

4.1 分片元数据管理

public class ShardingMetadataManager
{
    private readonly IRepository _configRepository;
    
    public async Task GetConfigAsync(string tableName)
    {
        return await _configRepository.GetByKeyAsync(tableName);
    }
    
    public async Task SaveConfigAsync(ShardingConfig config)
    {
        await _configRepository.SaveAsync(config);
    }
    
    public async Task> GetAllConfigsAsync()
    {
        return await _configRepository.GetAllAsync();
    }
    
    public async Task AddPartitionAsync(string tableName, RangePartition partition)
    {
        var config = await GetConfigAsync(tableName);
        
        if (config == null)
        {
            throw new NotFoundException("分片配置不存在");
        }
        
        config.Partitions.Add(partition);
        
        await SaveConfigAsync(config);
    }
    
    public async Task RemovePartitionAsync(string tableName, object startValue)
    {
        var config = await GetConfigAsync(tableName);
        
        if (config == null)
        {
            throw new NotFoundException("分片配置不存在");
        }
        
        config.Partitions.RemoveAll(p => Equals(p.StartValue, startValue));
        
        await SaveConfigAsync(config);
    }
    
    public async Task ValidateConfigAsync(ShardingConfig config)
    {
        if (string.IsNullOrEmpty(config.TableName))
        {
            return false;
        }
        
        if (config.StrategyType == ShardingStrategyType.Range && 
            (config.Partitions == null || config.Partitions.Count == 0))
        {
            return false;
        }
        
        return true;
    }
}

public class ShardingConfig
{
    public string TableName { get; set; }
    public ShardingStrategyType StrategyType { get; set; }
    public string ShardKey { get; set; }
    public int DatabaseCount { get; set; }
    public int TableCount { get; set; }
    public List Partitions { get; set; } = new();
    public DateTime CreatedAt { get; set; }
    public DateTime UpdatedAt { get; set; }
}

public enum ShardingStrategyType { Hash, Range, List, Composite }

4.2 分片迁移

public class ShardMigrationService
{
    public async Task MigrateDataAsync(string sourceDb, string sourceTable, 
        string targetDb, string targetTable, Func filter)
    {
        var result = new MigrationResult
        {
            SourceDatabase = sourceDb,
            SourceTable = sourceTable,
            TargetDatabase = targetDb,
            TargetTable = targetTable
        };
        
        try
        {
            var migratedCount = 0;
            var batchSize = 1000;
            
            await _dataMigrationService.StartMigrationAsync(sourceDb, sourceTable, 
                targetDb, targetTable, filter, batchSize);
            
            await _shardingMetadataManager.UpdateConfigAsync(targetDb, targetTable);
            
            await _cacheInvalidationService.BroadcastInvalidationAsync(sourceTable);
            
            result.Success = true;
            result.MigratedCount = migratedCount;
        }
        catch (Exception ex)
        {
            result.Success = false;
            result.ErrorMessage = ex.Message;
        }
        
        return result;
    }
    
    public async Task RebalanceAsync(string tableName)
    {
        var config = await _shardingMetadataManager.GetConfigAsync(tableName);
        
        if (config == null)
        {
            return new MigrationResult { Success = false, Message = "分片配置不存在" };
        }
        
        var dataDistribution = await _dataDistributionService.AnalyzeDistributionAsync(tableName);
        
        var migrations = CalculateRebalanceMigrations(dataDistribution, config);
        
        foreach (var migration in migrations)
        {
            await MigrateDataAsync(migration.SourceDb, migration.SourceTable, 
                migration.TargetDb, migration.TargetTable, migration.Filter);
        }
        
        return new MigrationResult { Success = true, Message = "数据重平衡完成" };
    }
    
    private List CalculateRebalanceMigrations(DataDistribution distribution, ShardingConfig config)
    {
        var migrations = new List();
        
        var averageCount = distribution.TotalCount / (config.DatabaseCount * config.TableCount);
        
        foreach (var partition in distribution.Partitions)
        {
            if (partition.RecordCount > averageCount * 1.2)
            {
                migrations.Add(new MigrationPlan
                {
                    SourceDb = partition.DatabaseName,
                    SourceTable = partition.TableName,
                    TargetDb = FindUnderutilizedPartition(distribution).DatabaseName,
                    TargetTable = FindUnderutilizedPartition(distribution).TableName,
                    Filter = BuildFilter(partition.ExcessRecords)
                });
            }
        }
        
        return migrations;
    }
}

public class MigrationResult
{
    public bool Success { get; set; }
    public string SourceDatabase { get; set; }
    public string SourceTable { get; set; }
    public string TargetDatabase { get; set; }
    public string TargetTable { get; set; }
    public int MigratedCount { get; set; }
    public string ErrorMessage { get; set; }
    public string Message { get; set; }
}

五、分布式事务

5.1 分布式事务方案

graph TD A[分布式事务] --> B{事务方案} B -->|2PC| C[两阶段提交] C --> C1[准备阶段] C1 --> C2[提交阶段] B -->|Saga| D[Saga模式] D --> D1[本地事务1] D1 --> D2[本地事务2] D2 --> D3[本地事务3] D3 --> D4[成功完成] D2 --> D5[补偿事务2] D5 --> D6[补偿事务1] B -->|TCC| E[TCC模式] E --> E1[Try] E1 --> E2[Confirm] E2 --> E3[成功完成] E1 --> E4[Cancel] B -->|消息队列| F[最终一致性] F --> F1[发送消息] F1 --> F2[本地事务] F2 --> F3[确认消息] F3 --> F4[消费者处理]

5.2 Saga模式实现

public class SagaTransactionManager
{
    private readonly List _steps = new();
    private readonly List _completedSteps = new();
    
    public void AddStep(ISagaStep step)
    {
        _steps.Add(step);
    }
    
    public async Task ExecuteAsync()
    {
        for (int i = 0; i < _steps.Count; i++)
        {
            try
            {
                await _steps[i].ExecuteAsync();
                _completedSteps.Add(_steps[i]);
            }
            catch (Exception ex)
            {
                await CompensateAsync();
                
                return new SagaResult
                {
                    Success = false,
                    FailedStep = i,
                    ErrorMessage = ex.Message
                };
            }
        }
        
        return new SagaResult { Success = true };
    }
    
    private async Task CompensateAsync()
    {
        for (int i = _completedSteps.Count - 1; i >= 0; i--)
        {
            try
            {
                await _completedSteps[i].CompensateAsync();
            }
            catch (Exception)
            {
            }
        }
    }
}

public interface ISagaStep
{
    Task ExecuteAsync();
    Task CompensateAsync();
}

public class SagaResult
{
    public bool Success { get; set; }
    public int FailedStep { get; set; }
    public string ErrorMessage { get; set; }
}

public class ShardedTransactionStep : ISagaStep
{
    private readonly string _databaseName;
    private readonly Func _operation;
    private readonly Func _compensation;
    
    public ShardedTransactionStep(string databaseName, Func operation, Func compensation)
    {
        _databaseName = databaseName;
        _operation = operation;
        _compensation = compensation;
    }
    
    public async Task ExecuteAsync()
    {
        using (_shardingRouter.SwitchToDatabase(_databaseName))
        {
            await _operation();
        }
    }
    
    public async Task CompensateAsync()
    {
        if (_compensation != null)
        {
            using (_shardingRouter.SwitchToDatabase(_databaseName))
            {
                await _compensation();
            }
        }
    }
}

六、分片最佳实践

6.1 分片设计原则

原则 描述 实现方式
选择合适的分片键 根据查询模式选择 用户ID、订单ID
避免跨分片查询 减少性能开销 合理设计表结构
考虑扩展性 支持动态扩缩容 一致性哈希
数据均匀分布 避免热点数据 哈希分片
事务处理 保证数据一致性 Saga模式

6.2 水平扩展最佳实践

  • 从单库开始,按需分片
  • 选择合适的分片策略
  • 考虑分片键的稳定性
  • 支持动态扩缩容
  • 监控分片数据分布

七、总结

数据库分片与水平扩展是处理海量数据的关键技术。通过选择合适的分片策略(范围分片、哈希分片、列表分片等),能够实现数据的均匀分布和高效查询。分片管理包括元数据管理和数据迁移,分布式事务处理是分片环境中必须考虑的问题。遵循分片设计原则,能够构建可扩展、高性能的分布式数据库系统。