1. 需求分析和评估阶段
1.1 业务需求调研
# 需求调研清单
requirements_checklist = {
"业务规模": {
"当前数据量": "评估现有数据规模",
"增长趋势": "预测未来1-3年增长",
"访问模式": "读写比例、峰值QPS"
},
"性能瓶颈": {
"慢查询": "分析慢查询日志",
"连接数": "数据库连接使用情况",
"存储空间": "磁盘使用率和增长趋势"
},
"业务特性": {
"数据关联": "表之间的关联关系",
"查询模式": "常用查询类型和频率",
"事务要求": "跨表事务需求"
}
}
1.2 技术可行性评估
— 评估当前数据库状态
SELECT
table_schema,
table_name,
table_rows,
ROUND(data_length/1024/1024, 2) AS data_mb,
ROUND(index_length/1024/1024, 2) AS index_mb,
ROUND((data_length + index_length)/1024/1024, 2) AS total_mb
FROM information_schema.tables
WHERE table_schema = 'your_database'
ORDER BY (data_length + index_length) DESC;
— 分析表访问频率
SELECT
table_name,
COUNT(*) as access_count
FROM information_schema.processlist
WHERE table_schema = 'your_database'
GROUP BY table_name
ORDER BY access_count DESC;
2. 方案设计阶段
2.1 分库分表策略选择
# 分库分表策略对比
sharding_strategies = {
"垂直分库": {
"适用场景": "业务模块清晰,表间关联少",
"优点": "实现简单,业务清晰",
"缺点": "可能存在跨库事务"
},
"水平分表": {
"适用场景": "单表数据量大,数据分布均匀",
"优点": "扩展性好,查询性能提升",
"缺点": "需要处理路由逻辑"
},
"混合方案": {
"适用场景": "复杂业务,需要综合考虑",
"优点": "兼顾业务和数据特性",
"缺点": "实现复杂度高"
}
}
# 分片键选择示例
sharding_keys = {
"用户表": "user_id", # 按用户ID分片
"订单表": "user_id", # 按用户ID分片,保证同一用户订单在同一分片
"商品表": "category_id", # 按分类分片
"日志表": "create_time" # 按时间分片
}
2.2 技术方案设计
// 分片规则配置示例
@Configuration
public class ShardingConfig {
@Bean
public DataSource shardingDataSource() throws SQLException {
ShardingRuleConfiguration shardingRuleConfig = new ShardingRuleConfiguration();
// 订单表分片配置
TableRuleConfiguration orderTableRule = new TableRuleConfiguration();
orderTableRule.setLogicTable("t_order");
orderTableRule.setActualDataNodes("ds_${0..1}.t_order_${0..3}");
// 分库策略
orderTableRule.setDatabaseShardingStrategyConfig(
new StandardShardingStrategyConfiguration("user_id",
new DatabaseShardingAlgorithm())
);
// 分表策略
orderTableRule.setTableShardingStrategyConfig(
new StandardShardingStrategyConfiguration("order_id",
new TableShardingAlgorithm())
);
shardingRuleConfig.getTableRuleConfigs().add(orderTableRule);
return ShardingDataSourceFactory.createDataSource(
createDataSourceMap(), shardingRuleConfig, new Properties()
);
}
}
// 分库算法
public class DatabaseShardingAlgorithm implements PreciseShardingAlgorithm<Long> {
@Override
public String doSharding(Collection<String> availableTargetNames,
PreciseShardingValue<Long> shardingValue) {
Long userId = shardingValue.getValue();
Long dbIndex = userId % 2;
return "ds_" + dbIndex;
}
}
3. 详细设计阶段
3.1 数据库设计
— 分库分表后的表结构设计
— 用户库0
CREATE DATABASE user_db_0;
USE user_db_0;
CREATE TABLE t_user_0 (
user_id BIGINT PRIMARY KEY,
username VARCHAR(50),
email VARCHAR(100),
create_time DATETIME,
update_time DATETIME,
INDEX idx_username (username),
INDEX idx_email (email)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
CREATE TABLE t_user_1 (
— 相同结构
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
— 订单库0
CREATE DATABASE order_db_0;
USE order_db_0;
CREATE TABLE t_order_0 (
order_id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
order_amount DECIMAL(10,2),
order_status TINYINT,
create_time DATETIME,
update_time DATETIME,
INDEX idx_user_id (user_id),
INDEX idx_create_time (create_time)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
3.2 应用层改造
# 数据访问层改造
class UserRepository:
def __init__(self, sharding_config):
self.sharding_config = sharding_config
self.databases = self._init_databases()
def get_user(self, user_id):
# 计算分片
db_index = user_id % self.sharding_config.db_count
table_index = user_id % self.sharding_config.table_count
# 路由到具体数据库和表
db = self.databases[f'user_db_{db_index}']
table = f't_user_{table_index}'
# 执行查询
query = f"SELECT * FROM {table} WHERE user_id = %s"
return db.execute(query, (user_id,))
def create_user(self, user_data):
user_id = self._generate_user_id()
db_index = user_id % self.sharding_config.db_count
table_index = user_id % self.sharding_config.table_count
db = self.databases[f'user_db_{db_index}']
table = f't_user_{table_index}'
query = f"INSERT INTO {table} (user_id, username, email) VALUES (%s, %s, %s)"
db.execute(query, (user_id, user_data['username'], user_data['email']))
return user_id
4. 开发实施阶段
4.1 数据迁移方案
# 数据迁移脚本
class DataMigrator:
def __init__(self, source_config, target_config):
self.source_db = self._connect(source_config)
self.target_dbs = self._connect_targets(target_config)
def migrate_users(self, batch_size=1000):
offset = 0
while True:
# 从源数据库读取数据
users = self.source_db.execute(
"SELECT * FROM t_user ORDER BY user_id LIMIT %s OFFSET %s",
(batch_size, offset)
)
if not users:
break
# 分片写入目标数据库
for user in users:
self._write_to_shard(user)
offset += batch_size
print(f"Migrated {offset} users")
def _write_to_shard(self, user):
user_id = user['user_id']
db_index = user_id % len(self.target_dbs)
table_index = user_id % 2 # 假设每个库2张表
target_db = self.target_dbs[db_index]
table = f't_user_{table_index}'
target_db.execute(
f"INSERT INTO {table} (user_id, username, email, create_time) VALUES (%s, %s, %s, %s)",
(user['user_id'], user['username'], user['email'], user['create_time'])
)
4.2 双写方案
// 双写服务实现
@Service
public class DualWriteService {
@Autowired
private OldUserRepository oldUserRepository;
@Autowired
private NewUserRepository newUserRepository;
@Autowired
private DataConsistencyChecker consistencyChecker;
@Transactional
public void createUser(User user) {
try {
// 先写新库
newUserRepository.save(user);
// 再写旧库
oldUserRepository.save(user);
// 记录操作日志
logOperation(user.getUserId(), "CREATE", "SUCCESS");
} catch (Exception e) {
// 处理失败情况
logOperation(user.getUserId(), "CREATE", "FAILED");
throw e;
}
}
@Scheduled(fixedDelay = 60000) // 每分钟检查一次
public void checkConsistency() {
List<Long> inconsistentUsers = consistencyChecker.findInconsistentUsers();
for (Long userId : inconsistentUsers) {
repairUserData(userId);
}
}
}
5. 测试验证阶段
5.1 功能测试
# 分库分表功能测试
class ShardingTest:
def test_user_crud(self):
# 测试创建用户
user_id = user_repository.create_user({
'username': 'test_user',
'email': 'test@example.com'
})
assert user_id is not None
# 测试查询用户
user = user_repository.get_user(user_id)
assert user['username'] == 'test_user'
# 测试更新用户
user_repository.update_user(user_id, {'email': 'new@example.com'})
updated_user = user_repository.get_user(user_id)
assert updated_user['email'] == 'new@example.com'
# 测试删除用户
user_repository.delete_user(user_id)
deleted_user = user_repository.get_user(user_id)
assert deleted_user is None
def test_sharding_routing(self):
# 测试分片路由正确性
user_ids = [1, 2, 3, 4, 5]
for user_id in user_ids:
user = user_repository.get_user(user_id)
expected_db = f'user_db_{user_id % 2}'
expected_table = f't_user_{user_id % 2}'
# 验证路由是否正确
assert user['db_name'] == expected_db
assert user['table_name'] == expected_table
5.2 性能测试
# 性能测试脚本
class PerformanceTest:
def test_concurrent_insert(self):
import threading
import time
def insert_users(thread_id, count):
start_time = time.time()
for i in range(count):
user_repository.create_user({
'username': f'user_{thread_id}_{i}',
'email': f'user_{thread_id}_{i}@example.com'
})
end_time = time.time()
print(f"Thread {thread_id} inserted {count} users in {end_time – start_time:.2f}s")
# 并发测试
threads = []
thread_count = 10
users_per_thread = 1000
start_time = time.time()
for i in range(thread_count):
thread = threading.Thread(target=insert_users, args=(i, users_per_thread))
threads.append(thread)
thread.start()
for thread in threads:
thread.join()
total_time = time.time() – start_time
total_users = thread_count * users_per_thread
print(f"Total: {total_users} users in {total_time:.2f}s")
print(f"Average: {total_users/total_time:.2f} users/second")
6. 上线部署阶段
6.1 灰度发布方案
# 灰度发布配置
gray_release:
phases:
– phase: 1
name: "内部测试"
traffic_percentage: 1
duration: "1天"
monitoring:
– error_rate
– response_time
– data_consistency
– phase: 2
name: "小流量验证"
traffic_percentage: 5
duration: "2天"
monitoring:
– error_rate
– response_time
– data_consistency
– system_load
– phase: 3
name: "中流量验证"
traffic_percentage: 20
duration: "3天"
monitoring:
– all_metrics
– phase: 4
name: "全量上线"
traffic_percentage: 100
duration: "持续"
monitoring:
– all_metrics
6.2 监控告警配置
# 监控指标配置
monitoring_metrics = {
"业务指标": {
"QPS": "每秒查询数",
"响应时间": "P50, P95, P99",
"错误率": "请求失败比例",
"数据一致性": "主从延迟、数据校验"
},
"系统指标": {
"CPU使用率": "数据库服务器CPU",
"内存使用率": "数据库服务器内存",
"磁盘IO": "读写IOPS",
"网络流量": "数据库网络带宽"
},
"分库分表特有指标": {
"分片分布": "各分片数据量分布",
"路由准确率": "路由正确性",
"跨片查询": "跨分片查询次数和耗时"
}
}
# 告警规则示例
alert_rules = {
"critical": {
"error_rate": "> 1%",
"response_time_p99": "> 1000ms",
"data_consistency": "存在不一致数据"
},
"warning": {
"error_rate": "> 0.1%",
"response_time_p99": "> 500ms",
"replication_delay": "> 5s"
}
}
7. 运维优化阶段
7.1 日常运维脚本
#!/bin/bash
# 分库分表运维脚本
# 检查各分片数据分布
check_sharding_distribution() {
echo "=== 分片数据分布检查 ==="
for db in user_db_0 user_db_1; do
echo "数据库: $db"
mysql -h $db_host -u $user -p$password $db -e "
SELECT
table_name,
table_rows,
ROUND(data_length/1024/1024, 2) AS data_mb
FROM information_schema.tables
WHERE table_schema = '$db'
ORDER BY table_rows;
"
done
}
# 数据一致性检查
check_data_consistency() {
echo "=== 数据一致性检查 ==="
# 实现数据对比逻辑
python3 check_consistency.py
}
# 性能监控
monitor_performance() {
echo "=== 性能监控 ==="
mysql -h $db_host -u $user -p$password -e "
SHOW ENGINE INNODB STATUS\\G
" | grep -A 10 "TRANSACTIONS"
}
# 主菜单
case "$1" in
distribution)
check_sharding_distribution
;;
consistency)
check_data_consistency
;;
performance)
monitor_performance
;;
*)
echo "Usage: $0 {distribution|consistency|performance}"
exit 1
esac
7.2 应急预案
# 应急处理预案
class EmergencyPlan:
def handle_shard_failure(self, shard_id):
"""处理分片故障"""
# 1. 将故障分片标记为不可用
self.mark_shard_unavailable(shard_id)
# 2. 将流量路由到备用分片
self.route_to_backup(shard_id)
# 3. 启动数据恢复
self.restore_shard(shard_id)
# 4. 恢复后重新加入集群
self.recover_shard(shard_id)
def handle_data_inconsistency(self, inconsistent_records):
"""处理数据不一致"""
for record in inconsistent_records:
# 1. 确定正确数据源
correct_data = self.determine_correct_data(record)
# 2. 修复不一致数据
self.repair_data(record, correct_data)
# 3. 记录修复日志
self.log_repair(record, correct_data)
def handle_traffic_spike(self):
"""处理流量激增"""
# 1. 启用限流
self.enable_rate_limiting()
# 2. 扩容分片
self.scale_out_shards()
# 3. 启用缓存
self.enable_cache()
8. 项目管理要点
8.1 时间规划
# 项目时间规划
project_timeline = {
"需求分析": "2周",
"方案设计": "3周",
"开发实施": "6周",
"测试验证": "4周",
"上线部署": "2周",
"运维优化": "持续"
}
# 里程碑设置
milestones = {
"M1": "需求分析完成",
"M2": "技术方案评审通过",
"M3": "核心功能开发完成",
"M4": "数据迁移完成",
"M5": "测试验证通过",
"M6": "灰度发布完成",
"M7": "全量上线"
}
8.2 风险管理
# 风险识别和应对
risk_management = {
"技术风险": {
"数据迁移失败": {
"概率": "中",
"影响": "高",
"应对措施": "充分测试,准备回滚方案"
},
"性能不达预期": {
"概率": "中",
"影响": "中",
"应对措施": "性能测试,优化查询,增加缓存"
}
},
"业务风险": {
"业务中断": {
"概率": "低",
"影响": "高",
"应对措施": "灰度发布,快速回滚机制"
},
"数据丢失": {
"概率": "低",
"影响": "高",
"应对措施": "数据备份,双写验证"
}
},
"项目风险": {
"进度延期": {
"概率": "中",
"影响": "中",
"应对措施": "合理排期,预留缓冲时间"
}
}
}
9. 团队协作要点
9.1 沟通机制
# 沟通计划
communication_plan = {
"日常沟通": {
"每日站会": "15分钟,同步进度和问题",
"技术讨论": "即时通讯工具",
"文档共享": "Wiki/文档中心"
},
"定期会议": {
"周会": "每周一次,总结本周工作",
"技术评审": "关键节点进行技术评审",
"项目汇报": "向管理层汇报项目进展"
},
"跨团队协作": {
"产品团队": "需求确认,业务逻辑验证",
"测试团队": "测试计划,用例评审",
"运维团队": "部署方案,监控告警",
"DBA团队": "数据库设计,性能优化"
}
}
9.2 文档管理
# 文档清单
documentation = {
"设计文档": [
"需求分析文档",
"技术方案设计",
"数据库设计文档",
"接口设计文档"
],
"开发文档": [
"开发规范",
"代码示例",
"API文档",
"配置说明"
],
"测试文档": [
"测试计划",
"测试用例",
"测试报告",
"性能测试报告"
],
"运维文档": [
"部署文档",
"运维手册",
"应急预案",
"监控告警说明"
]
}
10. 成功标准
10.1 技术指标
success_criteria = {
"性能指标": {
"查询响应时间": "P95 < 100ms",
"系统吞吐量": "支持当前业务3倍增长",
"数据库连接数": "连接池使用率 < 80%"
},
"可用性指标": {
"系统可用性": "> 99.9%",
"数据一致性": "无数据丢失",
"故障恢复时间": "< 5分钟"
},
"扩展性指标": {
"水平扩展": "支持在线增加分片",
"数据迁移": "支持在线数据迁移",
"负载均衡": "各分片负载均衡"
}
}
通过以上完整的实施流程,你可以系统性地主导分库分表项目。关键是要充分准备、严格测试、灰度上线,并建立完善的监控和应急机制。记住,分库分表是一个复杂的架构改造,需要谨慎规划和执行。

