一、分布式数据库概述
分布式数据库是将数据存储在多个物理节点上的数据库系统,通过分布式技术实现数据的水平扩展和高可用性。数据分片是分布式数据库的核心技术,通过将数据分散存储到多个节点来提高系统性能和容量。
二、分布式数据库对比
2.1 分布式数据库对比
| 数据库 | 架构 | 一致性 | 分片支持 | 适用场景 |
|---|---|---|---|---|
| TiDB | NewSQL | 强一致性 | 自动分片 | OLTP |
| CockroachDB | NewSQL | 强一致性 | 自动分片 | OLTP |
| OceanBase | NewSQL | 强一致性 | 自动分片 | 金融级 |
| ShardingSphere | 中间件 | 可配置 | 手动分片 | MySQL分片 |
| TDSQL | 云原生 | 强一致性 | 自动分片 | 云服务 |
2.2 分片策略对比
| 策略 | 特点 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|---|
| 哈希分片 | 按哈希值分布 | 均匀分布 | 数据均匀 | 范围查询差 |
| 范围分片 | 按范围分布 | 时间序列 | 范围查询好 | 热点问题 |
| 列表分片 | 按列表值分布 | 固定分类 | 灵活 | 维护复杂 |
| 复合分片 | 多键组合 | 复杂场景 | 灵活 | 复杂度高 |
三、分布式数据库架构
3.1 分布式数据库架构
graph TD
A[分布式数据库架构] --> B[Proxy层]
B --> C[分片策略]
C --> C1[哈希分片]
C1 --> C2[范围分片]
C2 --> C3[列表分片]
C3 --> C4[复合分片]
B --> D[数据节点]
D --> D1[节点1]
D --> D2[节点2]
D --> D3[节点3]
D --> D4[节点N]
D1 --> E1[分片1]
D1 --> E2[分片2]
D2 --> E3[分片3]
D2 --> E4[分片4]
D3 --> E5[分片5]
D3 --> E6[分片6]
F[元数据管理] --> F1[分片映射]
F1 --> F2[节点状态]
F2 --> F3[路由表]
G[一致性保障] --> G1[强一致性]
G1 --> G2[最终一致性]
G2 --> G3[分布式事务]
四、数据分片实现
4.1 分片策略实现
public class ShardingStrategy
{
public int GetShardIndex(object key, int shardCount)
{
var hash = key.GetHashCode();
return Math.Abs(hash) % shardCount;
}
public int GetShardIndexByHash(string key, int shardCount)
{
var hash = JenkinsHash(key);
return hash % shardCount;
}
public int GetShardIndexByRange(long value, List rangeBoundaries)
{
for (var i = 0; i < rangeBoundaries.Count; i++)
{
if (value <= rangeBoundaries[i])
{
return i;
}
}
return rangeBoundaries.Count;
}
public int GetShardIndexByList(string value, Dictionary listMapping)
{
return listMapping.TryGetValue(value, out var index) ? index : 0;
}
public int GetShardIndexByComposite(object[] keys, int shardCount)
{
var combinedKey = string.Join("|", keys);
return GetShardIndexByHash(combinedKey, shardCount);
}
private int JenkinsHash(string key)
{
var hash = 0;
foreach (var c in key)
{
hash += c;
hash += hash << 10;
hash ^= hash >> 6;
}
hash += hash << 3;
hash ^= hash >> 11;
hash += hash << 15;
return Math.Abs(hash);
}
}
public class ShardingConfiguration
{
public int ShardCount { get; set; }
public string ShardKey { get; set; }
public ShardingStrategyType StrategyType { get; set; }
public List RangeBoundaries { get; set; } = new();
public Dictionary ListMapping { get; set; } = new();
}
public enum ShardingStrategyType { Hash, Range, List, Composite }
4.2 分片路由
public class ShardingRouter
{
private readonly ShardingStrategy _strategy;
private readonly Dictionary _tableConfigs;
public ShardingRouter(ShardingStrategy strategy, Dictionary tableConfigs)
{
_strategy = strategy;
_tableConfigs = tableConfigs;
}
public ShardingResult Route(string tableName, object shardKey)
{
if (!_tableConfigs.TryGetValue(tableName, out var config))
{
throw new ShardingException($"Table {tableName} is not configured for sharding");
}
var shardIndex = config.StrategyType switch
{
ShardingStrategyType.Hash => _strategy.GetShardIndexByHash(shardKey.ToString(), config.ShardCount),
ShardingStrategyType.Range => _strategy.GetShardIndexByRange((long)shardKey, config.RangeBoundaries),
ShardingStrategyType.List => _strategy.GetShardIndexByList(shardKey.ToString(), config.ListMapping),
ShardingStrategyType.Composite => _strategy.GetShardIndexByComposite((object[])shardKey, config.ShardCount),
_ => throw new ShardingException("Unknown sharding strategy")
};
return new ShardingResult
{
TableName = $"{tableName}_{shardIndex}",
ShardIndex = shardIndex,
NodeIndex = shardIndex % config.ShardCount
};
}
public List RouteRange(string tableName, object startKey, object endKey)
{
if (!_tableConfigs.TryGetValue(tableName, out var config))
{
throw new ShardingException($"Table {tableName} is not configured for sharding");
}
var results = new List();
if (config.StrategyType == ShardingStrategyType.Range)
{
var startIndex = _strategy.GetShardIndexByRange((long)startKey, config.RangeBoundaries);
var endIndex = _strategy.GetShardIndexByRange((long)endKey, config.RangeBoundaries);
for (var i = startIndex; i <= endIndex; i++)
{
results.Add(new ShardingResult
{
TableName = $"{tableName}_{i}",
ShardIndex = i,
NodeIndex = i % config.ShardCount
});
}
}
return results;
}
}
public class ShardingResult
{
public string TableName { get; set; }
public int ShardIndex { get; set; }
public int NodeIndex { get; set; }
}
public class ShardingException : Exception
{
public ShardingException(string message) : base(message) { }
}
五、分布式事务
5.1 分布式事务实现
public class DistributedTransactionManager
{
private readonly ITransactionCoordinator _coordinator;
public DistributedTransactionManager(ITransactionCoordinator coordinator)
{
_coordinator = coordinator;
}
public async Task ExecuteTransactionAsync(string transactionId, Func> transactionalFunc)
{
await _coordinator.BeginTransactionAsync(transactionId);
try
{
var result = await transactionalFunc();
await _coordinator.CommitTransactionAsync(transactionId);
return result;
}
catch (Exception ex)
{
await _coordinator.RollbackTransactionAsync(transactionId);
throw new TransactionException("Transaction failed", ex);
}
}
public async Task ExecuteTransactionAsync(string transactionId, Func transactionalFunc)
{
await _coordinator.BeginTransactionAsync(transactionId);
try
{
await transactionalFunc();
await _coordinator.CommitTransactionAsync(transactionId);
}
catch (Exception ex)
{
await _coordinator.RollbackTransactionAsync(transactionId);
throw new TransactionException("Transaction failed", ex);
}
}
}
public interface ITransactionCoordinator
{
Task BeginTransactionAsync(string transactionId);
Task CommitTransactionAsync(string transactionId);
Task RollbackTransactionAsync(string transactionId);
}
public class SagaTransactionCoordinator : ITransactionCoordinator
{
private readonly Dictionary> _sagaSteps = new();
public async Task BeginTransactionAsync(string transactionId)
{
_sagaSteps[transactionId] = new List();
}
public async Task CommitTransactionAsync(string transactionId)
{
foreach (var step in _sagaSteps[transactionId])
{
await step.CompensateAsync();
}
}
public async Task RollbackTransactionAsync(string transactionId)
{
if (_sagaSteps.TryGetValue(transactionId, out var steps))
{
for (var i = steps.Count - 1; i >= 0; i--)
{
await steps[i].CompensateAsync();
}
}
}
public void RegisterStep(string transactionId, SagaStep step)
{
_sagaSteps[transactionId].Add(step);
}
}
public class SagaStep
{
public Func ExecuteAsync { get; set; }
public Func CompensateAsync { get; set; }
}
public class TransactionException : Exception
{
public TransactionException(string message, Exception innerException) : base(message, innerException) { }
}
六、分片扩容
6.1 分片扩容实现
public class ShardScaler
{
private readonly ShardingRouter _router;
private readonly IDatabaseMigration _migration;
public ShardScaler(ShardingRouter router, IDatabaseMigration migration)
{
_router = router;
_migration = migration;
}
public async Task ScaleUpAsync(string tableName, int newShardCount)
{
var config = GetTableConfig(tableName);
var oldShardCount = config.ShardCount;
for (var oldIndex = 0; oldIndex < oldShardCount; oldIndex++)
{
var sourceTable = $"{tableName}_{oldIndex}";
await _migration.CreateTableAsync(tableName, newShardCount);
await MigrateDataAsync(sourceTable, tableName, oldShardCount, newShardCount);
await _migration.DropTableAsync(sourceTable);
}
config.ShardCount = newShardCount;
}
private async Task MigrateDataAsync(string sourceTable, string targetTablePrefix, int oldCount, int newCount)
{
var data = await _migration.ReadDataAsync(sourceTable);
foreach (var record in data)
{
var shardKey = record[GetShardKey(targetTablePrefix)];
var newIndex = CalculateNewIndex(shardKey, newCount);
await _migration.InsertDataAsync($"{targetTablePrefix}_{newIndex}", record);
}
}
private int CalculateNewIndex(object shardKey, int newCount)
{
var hash = shardKey.GetHashCode();
return Math.Abs(hash) % newCount;
}
private ShardingConfiguration GetTableConfig(string tableName)
{
return new ShardingConfiguration { ShardCount = 4 };
}
}
public interface IDatabaseMigration
{
Task CreateTableAsync(string tableName, int shardCount);
Task DropTableAsync(string tableName);
Task>> ReadDataAsync(string tableName);
Task InsertDataAsync(string tableName, Dictionary record);
}
七、分布式数据库最佳实践
7.1 分片设计原则
| 原则 | 描述 | 实现方式 |
|---|---|---|
| 选择合适的分片键 | 均匀分布数据 | 用户ID/订单ID |
| 避免跨分片查询 | 减少性能开销 | 合理设计查询 |
| 考虑扩容 | 支持平滑扩容 | 预分片设计 |
| 数据一致性 | 保障数据一致性 | 分布式事务 |
| 监控告警 | 监控分片状态 | Prometheus |
7.2 分布式数据库最佳实践
- 选择合适的分片策略
- 合理设计分片键
- 实现分片路由
- 处理分布式事务
- 规划分片扩容
八、总结
分布式数据库与数据分片策略是构建大规模数据存储系统的核心技术。通过选择合适的分布式数据库和分片策略,能够实现数据的水平扩展和高可用性。分片路由、分布式事务和分片扩容是分布式数据库的关键实现。遵循分布式数据库最佳实践,能够构建高性能、高可用的数据存储系统。