数据库中间件从MyCAT到ShardingSphere的迁移复盘:分库分表策略的重构与性能压测对比
数据库中间件从MyCAT到ShardingSphere的迁移复盘:分库分表策略的重构与性能压测对比
一、项目背景与业务挑战
2024年底,我们运维团队接到一个令人头疼的迁移任务:将核心业务数据库的中间件从MyCAT迁移到ShardingSphere。这不是一次简单的"换个组件"操作,而是牵涉到分库分表策略重构、SQL兼容性适配、连接池管理重构、以及性能基准对比验证的系统工程。
MyCAT的维护困境。 MyCAT作为我们2018年选型的分库分表中间件,已经运行了6年。但随着业务规模增长,MyCAT的三个致命问题逐渐暴露:一是社区停止维护——最后一次正式Release是2019年,大量Bug和安全漏洞无人修复,我们自行维护了47个补丁但精力有限;二是SQL兼容性差——MyCAT只支持有限的SQL语法,复杂查询(子查询、跨分片JOIN、聚合函数)经常报错或返回错误结果,研发团队不得不将大量业务逻辑降级到应用层处理;三是性能瓶颈——MyCAT的单线程模型导致高并发下吞吐量受限,QPS超过8000时延迟急剧上升,P99延迟超过2秒。
ShardingSphere的优势与迁移的必要性。 Apache ShardingSphere作为活跃社区维护的分库分表中间件,在SQL兼容性、并发性能、生态集成三个方面远超MyCAT:支持99%的MySQL语法(包括子查询、跨分片JOIN、复杂聚合)、基于Netty的多线程NIO架构提供高并发吞吐、与Spring Boot/K8s/云原生生态深度集成。迁移到ShardingSphere是业务发展的刚性需求。
迁移的核心挑战。 这次迁移不是"停机替换"——业务不可中断。我们需要解决四个关键问题:一是分库分表策略差异——MyCAT使用rule.xml配置分片规则,ShardingSphere使用YAML配置,且分片算法的表达能力不同;二是SQL兼容性差异——大量在MyCAT中需要应用层 workaround的SQL,在ShardingSphere中可以直接执行,但也有一些MyCAT特有的SQL扩展需要适配;三是连接池管理差异——MyCAT自建连接池,ShardingSphere依赖HikariCP,连接数和超时配置需要重新调优;四是性能基准验证——必须通过压测证明ShardingSphere在吞吐量、延迟、资源消耗上不低于MyCAT,否则迁移不达标。
整个迁移周期从立项到全量切换历时4个月,涉及3个业务集群、12个逻辑库、48个分片、超过200TB数据。本文将完整复盘这次迁移的方案设计、实践落地、关键挑战和最终结果。
二、核心方案
2.1 四阶段渐进式迁移策略
我们设计了"评估-并行-验证-切换"四阶段迁移策略,每个阶段有明确的准入和退出条件:
2.2 分片策略映射与重构
MyCAT和ShardingSphere的分片配置存在结构性差异。MyCAT使用XML定义分片规则(tableRule + function),ShardingSphere使用YAML定义(shardingRule + shardingAlgorithms)。我们开发了自动化映射工具,将MyCAT的分片规则解析并转化为ShardingSphere的YAML配置。
关键映射逻辑:
- MyCAT的
PartitionByMod算法 → ShardingSphere的MOD分片算法 - MyCAT的
PartitionByLong(范围分片) → ShardingSphere的RANGE_INTERVAL分片算法 - MyCAT的
PartitionByHashMod→ ShardingSphere的HASH_MOD分片算法 - MyCAT的
PartitionByDate(按日期分片) → ShardingSphere的INTERVAL分片算法
但不是所有映射都是一一对应的——MyCAT支持自定义Java分片算法,ShardingSphere也支持但接口不同。对于我们自行开发的3个自定义分片算法(按业务ID前缀分片、按地域分片、按用户等级分片),需要按照ShardingSphere的ShardingAlgorithm接口重新实现。
2.3 性能压测对比设计
性能压测是迁移决策的关键依据。我们搭建了两组对照集群:MyCAT集群(2 Proxy实例)和ShardingSphere集群(2 Proxy实例),后端连接相同的物理MySQL分片,确保压测对比的公平性。压测维度包括:吞吐量(QPS)、延迟(P50/P95/P99)、连接池效率、SQL兼容性通过率、资源消耗(CPU/Memory/Network)、跨分片查询性能。
三、实践落地
3.1 分片策略映射工具
from dataclasses import dataclass, field
from typing import List, Dict, Optional, Tuple
from enum import Enum
import xml.etree.ElementTree as ET
import yaml
import re
import logging
logger = logging.getLogger(__name__)
class MyCATAlgorithmType(Enum):
"""MyCAT分片算法类型"""
MOD = "PartitionByMod" # 取模分片
LONG_RANGE = "PartitionByLong" # 范围分片
HASH_MOD = "PartitionByHashMod" # 哈希取模分片
DATE = "PartitionByDate" # 按日期分片
CUSTOM = "CustomAlgorithm" # 自定义算法
class ShardingSphereAlgorithmType(Enum):
"""ShardingSphere分片算法类型"""
MOD = "MOD" # 取模分片
RANGE_INTERVAL = "RANGE_INTERVAL" # 范围分片
HASH_MOD = "HASH_MOD" # 哈希取模分片
INTERVAL = "INTERVAL" # 按日期分片
CLASS_BASED = "CLASS_BASED" # 基于类的自定义算法
# MyCAT → ShardingSphere 算法类型映射
ALGORITHM_TYPE_MAPPING = {
MyCATAlgorithmType.MOD: ShardingSphereAlgorithmType.MOD,
MyCATAlgorithmType.LONG_RANGE: ShardingSphereAlgorithmType.RANGE_INTERVAL,
MyCATAlgorithmType.HASH_MOD: ShardingSphereAlgorithmType.HASH_MOD,
MyCATAlgorithmType.DATE: ShardingSphereAlgorithmType.INTERVAL,
MyCATAlgorithmType.CUSTOM: ShardingSphereAlgorithmType.CLASS_BASED,
}
@dataclass
class MyCATTableRule:
"""MyCAT分片规则"""
table_name: str # 逻辑表名
data_node: str # 数据节点分布
algorithm_type: MyCATAlgorithmType # 分片算法类型
algorithm_params: Dict[str, str] # 算法参数
partition_count: int # 分片数量
primary_key: str # 分片键
@dataclass
class ShardingSphereRule:
"""ShardingSphere分片规则"""
logic_table: str # 逻辑表名
actual_data_nodes: str # 实际数据节点分布表达式
sharding_column: str # 分片列
sharding_algorithm_name: str # 分片算法名称
sharding_algorithm_type: ShardingSphereAlgorithmType # 分片算法类型
sharding_algorithm_props: Dict[str, str] # 算法属性
@dataclass
class MigrationMappingResult:
"""迁移映射结果"""
logic_schema: str # 逻辑库名
mycat_rules: List[MyCATTableRule] # MyCAT原始规则
sharding_sphere_rules: List[ShardingSphereRule] # 映射后的ShardingSphere规则
unmapped_rules: List[MyCATTableRule] # 无法自动映射的规则
compatibility_issues: List[Dict] # 兼容性问题列表
class MyCATToShardingSphereMapper:
"""MyCAT分片规则 → ShardingSphere YAML配置的自动映射工具"""
def __init__(self, mycat_config_path: str, cluster_config: Dict):
self.mycat_config_path = mycat_config_path
self.cluster_config = cluster_config # 集群配置(分片命名规则等)
def map_rules(self) -> MigrationMappingResult:
"""执行分片规则映射"""
try:
# Step 1: 解析MyCAT rule.xml
mycat_rules = self._parse_mycat_rule_xml()
# Step 2: 逐规则映射到ShardingSphere格式
ss_rules = []
unmapped = []
issues = []
for rule in mycat_rules:
try:
ss_rule = self._map_single_rule(rule)
ss_rules.append(ss_rule)
except Exception as e:
logger.warning(f"规则映射失败: {rule.table_name}, 原因: {e}")
unmapped.append(rule)
issues.append({
"table": rule.table_name,
"reason": str(e),
"type": "algorithm_mapping_failure",
})
# Step 3: 生成ShardingSphere YAML配置
result = MigrationMappingResult(
logic_schema=self.cluster_config.get("schema_name", "ds"),
mycat_rules=mycat_rules,
sharding_sphere_rules=ss_rules,
unmapped_rules=unmapped,
compatibility_issues=issues,
)
logger.info(
f"规则映射完成: 总规则={len(mycat_rules)}, "
f"成功映射={len(ss_rules)}, 未映射={len(unmapped)}, "
f"兼容性问题={len(issues)}"
)
return result
except Exception as e:
logger.error(f"规则映射异常: {e}", exc_info=True)
raise
def _parse_mycat_rule_xml(self) -> List[MyCATTableRule]:
"""解析MyCAT rule.xml文件"""
rules = []
try:
tree = ET.parse(self.mycat_config_path)
root = tree.getroot()
# 解析tableRule节点
for table_rule in root.findall(".//tableRule"):
rule_name = table_rule.get("name", "")
rule_node = table_rule.find("rule")
if rule_node is None:
continue
columns = rule_node.find("columns").text.strip() if rule_node.find("columns") is not None else ""
algorithm_name = rule_node.find("algorithm").text.strip() if rule_node.find("algorithm") is not None else ""
# 解析对应的function节点获取算法参数
func_node = root.find(f".//function[@name='{algorithm_name}']")
if func_node is None:
logger.warning(f"未找到算法定义: {algorithm_name}")
continue
algorithm_class = func_node.get("class", "")
algorithm_type = self._identify_algorithm_type(algorithm_class)
algorithm_params = self._extract_algorithm_params(func_node)
partition_count = int(func_node.get("partitionCount", "0"))
rules.append(MyCATTableRule(
table_name=rule_name,
data_node=self._resolve_data_nodes(rule_name),
algorithm_type=algorithm_type,
algorithm_params=algorithm_params,
partition_count=partition_count,
primary_key=columns,
))
except ET.ParseError as e:
logger.error(f"MyCAT rule.xml解析失败: {e}")
raise
return rules
def _identify_algorithm_type(self, class_name: str) -> MyCATAlgorithmType:
"""识别MyCAT算法类型"""
# 常见MyCAT分片算法的类名映射
CLASS_TO_TYPE = {
"io.mycat.route.function.PartitionByMod": MyCATAlgorithmType.MOD,
"io.mycat.route.function.PartitionByLong": MyCATAlgorithmType.LONG_RANGE,
"io.mycat.route.function.PartitionByHashMod": MyCATAlgorithmType.HASH_MOD,
"io.mycat.route.function.PartitionByDate": MyCATAlgorithmType.DATE,
}
for cls, alg_type in CLASS_TO_TYPE.items():
if cls in class_name:
return alg_type
# 自定义算法
return MyCATAlgorithmType.CUSTOM
def _extract_algorithm_params(self, func_node) -> Dict[str, str]:
"""提取MyCAT算法参数"""
params = {}
for attr_name in ["partitionCount", "partitionLength", "mapFile",
"beginDate", "endDate", "partitions", "defaultNode"]:
attr_value = func_node.get(attr_name)
if attr_value:
params[attr_name] = attr_value
# 提取property子节点参数
for prop in func_node.findall("property"):
prop_name = prop.get("name", "")
prop_value = prop.text.strip() if prop.text else ""
if prop_name and prop_value:
params[prop_name] = prop_value
return params
def _map_single_rule(self, mycat_rule: MyCATTableRule) -> ShardingSphereRule:
"""映射单条MyCAT规则到ShardingSphere格式"""
# 确定目标算法类型
ss_algorithm_type = ALGORITHM_TYPE_MAPPING.get(mycat_rule.algorithm_type)
if ss_algorithm_type is None:
raise ValueError(f"无法映射算法类型: {mycat_rule.algorithm_type}")
# 确定算法名称
algorithm_name = f"{mycat_rule.table_name}_sharding_algorithm"
# 映射算法参数
ss_props = self._map_algorithm_props(
mycat_rule.algorithm_type, ss_algorithm_type,
mycat_rule.algorithm_params, mycat_rule.partition_count
)
# 映射数据节点分布表达式
actual_data_nodes = self._map_data_nodes(
mycat_rule.table_name, mycat_rule.data_node, mycat_rule.partition_count
)
return ShardingSphereRule(
logic_table=mycat_rule.table_name,
actual_data_nodes=actual_data_nodes,
sharding_column=mycat_rule.primary_key,
sharding_algorithm_name=algorithm_name,
sharding_algorithm_type=ss_algorithm_type,
sharding_algorithm_props=ss_props,
)
def _map_algorithm_props(self, src_type: MyCATAlgorithmType,
target_type: ShardingSphereAlgorithmType,
params: Dict[str, str],
partition_count: int) -> Dict[str, str]:
"""映射算法参数"""
if target_type == ShardingSphereAlgorithmType.MOD:
return {
"sharding-count": str(partition_count),
}
elif target_type == ShardingSphereAlgorithmType.RANGE_INTERVAL:
# MyCAT范围分片的mapFile → ShardingSphere的range-lower/cupper
map_file_content = params.get("mapFile_content", "")
return self._parse_range_map_file(map_file_content, partition_count)
elif target_type == ShardingSphereAlgorithmType.HASH_MOD:
return {
"sharding-count": str(partition_count),
}
elif target_type == ShardingSphereAlgorithmType.INTERVAL:
return {
"datetime-pattern": params.get("dateFormat", "yyyy-MM-dd"),
"datetime-interval-amount": params.get("step", "1"),
"datetime-interval-unit": params.get("stepUnit", "MONTHS"),
"sharding-scope-pattern-offset": "0",
}
elif target_type == ShardingSphereAlgorithmType.CLASS_BASED:
# 自定义算法需要指定实现类
custom_class = params.get("customClass", "")
return {
"strategy": "STANDARD",
"algorithmClassName": self._adapt_custom_class(custom_class),
}
else:
raise ValueError(f"未知目标算法类型: {target_type}")
def _map_data_nodes(self, table_name: str, data_node: str,
partition_count: int) -> str:
"""映射数据节点分布表达式"""
# MyCAT: dn0,dn1,dn2,... → ShardingSphere: ds_${0..partition_count-1}.table_name
naming_rule = self.cluster_config.get("node_naming_rule", "ds_${0..%d}")
ds_pattern = naming_rule % (partition_count - 1)
return f"{ds_pattern}.{table_name}"
def generate_yaml_config(self, mapping_result: MigrationMappingResult) -> str:
"""生成ShardingSphere YAML配置文件"""
config = {
"dataSources": self._generate_datasource_config(),
"rules": [{
"shardingRule": {
"tables": {},
"shardingAlgorithms": {},
"bindingTables": [],
}
}]
}
# 填充分片表配置
for rule in mapping_result.sharding_sphere_rules:
config["rules"][0]["shardingRule"]["tables"][rule.logic_table] = {
"actualDataNodes": rule.actual_data_nodes,
"shardingColumn": rule.sharding_column,
"shardingAlgorithmName": rule.sharding_algorithm_name,
}
# 填充分片算法配置
config["rules"][0]["shardingRule"]["shardingAlgorithms"][rule.sharding_algorithm_name] = {
"type": rule.sharding_algorithm_type.value,
"props": rule.sharding_algorithm_props,
}
return yaml.dump(config, allow_unicode=True, default_flow_style=False, sort_keys=False)
def _generate_datasource_config(self) -> Dict:
"""生成数据源配置"""
datasources = {}
for i in range(self.cluster_config.get("shard_count", 0)):
ds_name = f"ds_{i}"
datasources[ds_name] = {
"url": f"jdbc:mysql://{self.cluster_config.get('mysql_host')}:3306/shard_{i}",
"username": self.cluster_config.get("mysql_user", "root"),
"password": self.cluster_config.get("mysql_password", ""),
"connectionTimeoutMilliseconds": 30000,
"idleTimeoutMilliseconds": 60000,
"maxLifetimeMilliseconds": 1800000,
"maxPoolSize": 50,
"minPoolSize": 1,
}
return datasources
def _adapt_custom_class(self, mycat_class: str) -> str:
"""适配自定义分片算法的类名(MyCAT接口→ShardingSphere接口)"""
# 记录需要重写的自定义算法
logger.info(f"自定义算法需要重写: {mycat_class}")
return f"com.sharding.custom.{mycat_class.split('.')[-1]}"
3.2 SQL兼容性扫描器
from dataclasses import dataclass
from typing import List, Dict, Tuple
import logging
logger = logging.getLogger(__name__)
@dataclass
class SQLCompatibilityResult:
"""SQL兼容性测试结果"""
sql_id: str
original_sql: str
mycat_result: str # MyCAT执行结果描述
ss_result: str # ShardingSphere执行结果描述
is_compatible: bool # 是否兼容
difference_type: str # 差异类型:syntax_error / result_diff / performance_diff / ok
class SQLCompatibilityScanner:
"""SQL兼容性扫描器:业务SQL在ShardingSphere上的执行测试"""
def __init__(self, mycat_executor, ss_executor, sql_repository):
self.mycat_executor = mycat_executor # MyCAT执行器
self.ss_executor = ss_executor # ShardingSphere执行器
self.sql_repository = sql_repository # SQL仓库(业务历史SQL)
def scan(self, schema_name: str) -> Dict:
"""扫描指定逻辑库的所有业务SQL兼容性"""
try:
# Step 1: 从SQL仓库获取该逻辑库的历史SQL
sql_list = self.sql_repository.get_sqls_by_schema(schema_name)
logger.info(f"获取到 {schema_name} 的历史SQL: {len(sql_list)}条")
results = []
compatible_count = 0
incompatible_count = 0
diff_types = {}
# Step 2: 逐条SQL在两个中间件上执行对比
for sql_info in sql_list:
try:
result = self._test_single_sql(sql_info)
results.append(result)
if result.is_compatible:
compatible_count += 1
else:
incompatible_count += 1
diff_types[result.difference_type] = diff_types.get(result.difference_type, 0) + 1
except Exception as e:
logger.warning(f"SQL测试异常: {sql_info['id']}, {e}")
incompatible_count += 1
diff_types["test_error"] = diff_types.get("test_error", 0) + 1
# Step 3: 计算兼容率并生成报告
compatibility_rate = compatible_count / len(sql_list) if sql_list else 0
report = {
"schema": schema_name,
"total_sqls": len(sql_list),
"compatible_count": compatible_count,
"incompatible_count": incompatible_count,
"compatibility_rate": compatibility_rate,
"diff_type_distribution": diff_types,
"incompatible_sqls": [r for r in results if not r.is_compatible],
}
logger.info(
f"SQL兼容性扫描完成: {schema_name}, "
f"兼容率={compatibility_rate:.2%}, "
f"不兼容={incompatible_count}条"
)
return report
except Exception as e:
logger.error(f"SQL兼容性扫描异常: {e}", exc_info=True)
return {"schema": schema_name, "error": str(e)}
def _test_single_sql(self, sql_info: Dict) -> SQLCompatibilityResult:
"""测试单条SQL在两个中间件上的执行结果"""
sql_id = sql_info["id"]
sql_text = sql_info["sql"]
# 在MyCAT上执行(对照组)
try:
mycat_result = self.mycat_executor.execute(
sql=sql_text, schema=sql_info.get("schema")
)
except Exception as e:
mycat_result = f"MyCAT执行错误: {e}"
# 在ShardingSphere上执行(实验组)
try:
ss_result = self.ss_executor.execute(
sql=sql_text, schema=sql_info.get("schema")
)
except Exception as e:
ss_result = f"ShardingSphere执行错误: {e}"
# 对比结果
is_compatible, diff_type = self._compare_results(mycat_result, ss_result, sql_text)
return SQLCompatibilityResult(
sql_id=sql_id,
original_sql=sql_text,
mycat_result=str(mycat_result)[:200],
ss_result=str(ss_result)[:200],
is_compatible=is_compatible,
difference_type=diff_type,
)
def _compare_results(self, mycat_result, ss_result, sql_text: str) -> Tuple[bool, str]:
"""对比两个中间件的执行结果"""
# 两者都执行错误
if isinstance(mycat_result, str) and "错误" in mycat_result:
if isinstance(ss_result, str) and "错误" in ss_result:
# 同样不支持 → 兼容(行为一致)
return True, "ok"
# ShardingSphere执行成功但MyCAT不支持 → 兼容性提升
if isinstance(mycat_result, str) and "错误" in mycat_result:
if not isinstance(ss_result, str) or "错误" not in ss_result:
return True, "ok" # ShardingSphere支持了MyCAT不支持的SQL
# ShardingSphere执行错误但MyCAT支持 → 不兼容
if isinstance(ss_result, str) and "错误" in ss_result:
if not isinstance(mycat_result, str) or "错误" not in mycat_result:
return False, "syntax_error"
# 两者都成功,对比结果集
if isinstance(mycat_result, list) and isinstance(ss_result, list):
if len(mycat_result) == len(ss_result):
# 行数一致,检查内容差异
if mycat_result == ss_result:
return True, "ok"
else:
return False, "result_diff"
else:
return False, "result_diff"
# 执行时间差异(性能维度)
# 这里简化处理,实际需要更细粒度的性能对比
return True, "ok"
3.3 性能压测对比框架
from dataclasses import dataclass, field
from typing import List, Dict
import time
import statistics
import logging
logger = logging.getLogger(__name__)
@dataclass
class BenchmarkConfig:
"""压测配置"""
scenario_name: str
qps_target: int # 目标QPS
duration_seconds: int # 持续时间
concurrency: int # 并发数
sql_mix_ratio: Dict[str, float] # SQL类型占比:point_select/range_scan/join/update/insert
@dataclass
class BenchmarkMetric:
"""压测指标"""
scenario: str
middleware: str # mycat / sharding_sphere
actual_qps: float # 实际QPS
latency_p50_ms: float # P50延迟
latency_p95_ms: float # P95延迟
latency_p99_ms: float # P99延迟
error_rate: float # 错误率
cpu_usage_percent: float # CPU使用率
memory_usage_mb: float # 内存使用
connection_pool_active: int # 连接池活跃连接数
connection_pool_idle: int # 连接池空闲连接数
cross_shard_query_ratio: float # 跨分片查询占比
@dataclass
class BenchmarkComparison:
"""压测对比结果"""
scenario: str
mycat_metrics: BenchmarkMetric
ss_metrics: BenchmarkMetric
qps_improvement: float # QPS提升百分比
latency_p99_improvement: float # P99延迟改善百分比
error_rate_improvement: float # 错误率改善百分比
is_ss_better: bool # ShardingSphere是否整体优于MyCAT
# 压测场景配置
BENCHMARK_SCENARIOS = {
"point_select": BenchmarkConfig(
scenario_name="单分片点查",
qps_target=10000,
duration_seconds=300,
concurrency=100,
sql_mix_ratio={"point_select": 1.0},
),
"mixed workload": BenchmarkConfig(
scenario_name="混合读写负载",
qps_target=8000,
duration_seconds=300,
concurrency=80,
sql_mix_ratio={
"point_select": 0.45,
"range_scan": 0.15,
"join": 0.10,
"update": 0.20,
"insert": 0.10,
},
),
"cross_shard_join": BenchmarkConfig(
scenario_name="跨分片JOIN查询",
qps_target=2000,
duration_seconds=300,
concurrency=20,
sql_mix_ratio={"join": 1.0},
),
"high_concurrency_write": BenchmarkConfig(
scenario_name="高并发写入",
qps_target=5000,
duration_seconds=300,
concurrency=50,
sql_mix_ratio={"insert": 0.6, "update": 0.4},
),
}
class BenchmarkRunner:
"""性能压测对比框架"""
def __init__(self, mycat_proxy, ss_proxy, metrics_collector):
self.mycat_proxy = mycat_proxy # MyCAT Proxy实例
self.ss_proxy = ss_proxy # ShardingSphere Proxy实例
self.metrics_collector = metrics_collector # 指标采集器
def run_comparison(self) -> List[BenchmarkComparison]:
"""运行全部压测场景对比"""
comparisons = []
for scenario_name, config in BENCHMARK_SCENARIOS.items():
try:
logger.info(f"开始压测场景: {scenario_name}")
# 在MyCAT上执行压测
mycat_metrics = self._run_benchmark("mycat", config)
# 在ShardingSphere上执行压测
ss_metrics = self._run_benchmark("sharding_sphere", config)
# 对比结果
comparison = self._compare_metrics(scenario_name, mycat_metrics, ss_metrics)
comparisons.append(comparison)
logger.info(
f"压测对比完成: {scenario_name}, "
f"ShardingSphere整体优于MyCAT={comparison.is_ss_better}, "
f"QPS提升={comparison.qps_improvement:.1f}%, "
f"P99延迟改善={comparison.latency_p99_improvement:.1f}%"
)
except Exception as e:
logger.error(f"压测场景 {scenario_name} 异常: {e}", exc_info=True)
return comparisons
def _run_benchmark(self, middleware: str, config: BenchmarkConfig) -> BenchmarkMetric:
"""在指定中间件上运行压测"""
proxy = self.mycat_proxy if middleware == "mycat" else self.ss_proxy
# 采集压测前的资源基线
baseline = self.metrics_collector.collect_resource_baseline(middleware)
# 执行压测
latencies = []
error_count = 0
total_count = 0
start_time = time.time()
while time.time() - start_time < config.duration_seconds:
# 根据SQL类型占比选择本次执行的SQL
sql_type = self._pick_sql_type(config.sql_mix_ratio)
sql = self._generate_test_sql(sql_type, config)
try:
exec_start = time.time()
result = proxy.execute(sql, timeout=10)
exec_latency = (time.time() - exec_start) * 1000 # 转为毫秒
latencies.append(exec_latency)
if result.get("error"):
error_count += 1
except Exception as e:
error_count += 1
logger.debug(f"压测SQL执行异常: {e}")
total_count += 1
# 计算压测指标
actual_qps = total_count / config.duration_seconds
error_rate = error_count / total_count if total_count > 0 else 1.0
sorted_latencies = sorted(latencies)
p50 = sorted_latencies[int(len(sorted_latencies) * 0.50)] if sorted_latencies else 0
p95 = sorted_latencies[int(len(sorted_latencies) * 0.95)] if sorted_latencies else 0
p99 = sorted_latencies[int(len(sorted_latencies) * 0.99)] if sorted_latencies else 0
# 采集压测后的资源指标
post_resources = self.metrics_collector.collect_resource_snapshot(middleware)
return BenchmarkMetric(
scenario=config.scenario_name,
middleware=middleware,
actual_qps=actual_qps,
latency_p50_ms=p50,
latency_p95_ms=p95,
latency_p99_ms=p99,
error_rate=error_rate,
cpu_usage_percent=post_resources.get("cpu_usage", 0),
memory_usage_mb=post_resources.get("memory_usage_mb", 0),
connection_pool_active=post_resources.get("active_connections", 0),
connection_pool_idle=post_resources.get("idle_connections", 0),
cross_shard_query_ratio=post_resources.get("cross_shard_ratio", 0),
)
def _pick_sql_type(self, mix_ratio: Dict[str, float]) -> str:
"""根据SQL类型占比随机选择本次执行的SQL类型"""
import random
types = list(mix_ratio.keys())
weights = list(mix_ratio.values())
return random.choices(types, weights=weights, k=1)[0]
def _generate_test_sql(self, sql_type: str, config: BenchmarkConfig) -> str:
"""生成测试SQL"""
SQL_TEMPLATES = {
"point_select": "SELECT * FROM t_order WHERE order_id = {id}",
"range_scan": "SELECT * FROM t_order WHERE order_id BETWEEN {id} AND {id}+100",
"join": "SELECT o.*, u.name FROM t_order o JOIN t_user u ON o.user_id = u.id WHERE o.order_id = {id}",
"update": "UPDATE t_order SET status = 1 WHERE order_id = {id}",
"insert": "INSERT INTO t_order (order_id, user_id, status) VALUES ({id}, {id}%1000, 0)",
}
template = SQL_TEMPLATES.get(sql_type, SQL_TEMPLATES["point_select"])
return template.replace("{id}", str(int(time.time() * 1000) % 100000))
def _compare_metrics(self, scenario: str,
mycat: BenchmarkMetric,
ss: BenchmarkMetric) -> BenchmarkComparison:
"""对比两个中间件的压测指标"""
# QPS提升百分比
qps_improvement = (
(ss.actual_qps - mycat.actual_qps) / mycat.actual_qps * 100
if mycat.actual_qps > 0 else 0
)
# P99延迟改善百分比(减少为改善)
latency_p99_improvement = (
(mycat.latency_p99_ms - ss.latency_p99_ms) / mycat.latency_p99_ms * 100
if mycat.latency_p99_ms > 0 else 0
)
# 错误率改善百分比
error_rate_improvement = (
(mycat.error_rate - ss.error_rate) / mycat.error_rate * 100
if mycat.error_rate > 0 else 0
)
# 综合判定:QPS提升+P99改善+错误率改善 三项均为正 → ShardingSphere优于MyCAT
is_better = (
qps_improvement > 0 and
latency_p99_improvement > 0 and
error_rate_improvement >= 0
)
return BenchmarkComparison(
scenario=scenario,
mycat_metrics=mycat,
ss_metrics=ss,
qps_improvement=qps_improvement,
latency_p99_improvement=latency_p99_improvement,
error_rate_improvement=error_rate_improvement,
is_ss_better=is_better,
)
四、关键挑战与应对策略
4.1 自定义分片算法的重写与验证
MyCAT有3个自定义分片算法(按业务ID前缀分片、按地域分片、按用户等级分片),ShardingSphere的ShardingAlgorithm接口与MyCAT的分片接口差异显著——MyCAT返回Integer分片序号,ShardingSphere返回Collection<String>数据节点名。重写过程需要:一是将分片结果从整数索引映射到数据节点名格式(ds_0而非0);二是ShardingSphere的分片范围查询需要同时返回路由到的多个数据节点,而MyCAT只返回一个。我们在重写后对3个算法分别执行了10万次路由验证(覆盖全值域),确保路由结果与原算法100%一致。
4.2 跨分片JOIN的性能差异
MyCAT对跨分片JOIN采用"笛卡尔积+内存合并"策略,性能极差(2表JOIN超过3秒)。ShardingSphere采用"流式归并+内存限制"策略,同样场景下JOIN耗时约1.2秒,但仍远慢于单分片查询。对于核心业务的4个高频跨分片JOIN查询,我们设计了"宽表预聚合"方案——将JOIN结果预先聚合到宽表中,查询时直接读取宽表而非实时JOIN,4个查询的P99延迟从1.2秒降到50毫秒。
4.3 连接池参数的重新调优
MyCAT自建连接池默认配置(maxActive=100, idleCheck=60s),ShardingSphere使用HikariCP(默认maxPoolSize=50, idleTimeout=600s)。初期我们直接使用HikariCP默认配置,结果在高并发场景下出现连接等待超时(等待超过30秒)。通过4轮参数调优(基于压测反馈逐步调整),最终将maxPoolSize从50调整到80、connectionTimeout从30s调整到5s、idleTimeout从600s调整到300s、maxLifetime从1800s调整到1200s。调优后连接等待超时消失,连接池利用率从45%提升到72%。
4.4 双写期间的数据一致性保障
阶段二的"双中间件并行运行"需要确保写入数据在两个中间件中一致。我们采用"主写MyCAT+同步复制到ShardingSphere"的策略——应用层写入流量仍走MyCAT(保持业务稳定),同时由同步服务将写入事件复制到ShardingSphere。数据一致性校验采用"逐分片CRC校验"——每10分钟对每个分片的表数据计算CRC32摘要并对比,不一致时分片标记为"需修复"并自动触发修复流程。双写期间发现3次数据不一致(均因并发写入时的时序差异),均在10分钟内自动修复。
4.5 MyCAT特有SQL的适配问题
MyCAT有一些特有的SQL扩展语法(如/*!mycat:sql=select * from ...*/注释路由提示),这些语法在ShardingSphere中不支持。我们排查出业务代码中有23处使用了MyCAT特有语法,需要逐一手动适配。适配策略:一是将路由提示替换为ShardingSphere的HintManager API(Java代码中设置分片路由Hint);二是将MyCAT注释语法中的SQL提取出来直接执行(去掉注释包装)。23处适配在2周内完成,每处适配后都经过回归测试验证。
五、总结
MyCAT到ShardingSphere的迁移是一次深度基础设施替换项目,4个月的渐进式迁移最终取得了以下成果:
性能对比数据(4个压测场景汇总):
| 压测场景 | MyCAT QPS | SS QPS | QPS提升 | MyCAT P99(ms) | SS P99(ms) | P99改善 |
|---|---|---|---|---|---|---|
| 单分片点查 | 8,200 | 14,500 | +77% | 120 | 45 | +62% |
| 混合读写 | 6,500 | 10,800 | +66% | 1,850 | 380 | +79% |
| 跨分片JOIN | 1,200 | 2,800 | +133% | 3,200 | 1,200 | +62% |
| 高并发写入 | 4,800 | 7,600 | +58% | 350 | 180 | +49% |
- SQL兼容性:MyCAT通过率72%(28%的SQL不支持),ShardingSphere通过率96%(4%的SQL需适配),兼容性提升24个百分点
- 迁移成功率:48个分片全部成功迁移,0数据丢失,0服务中断
- 回滚预案:双写阶段保留MyCAT作为回退通道,全量切换后保留2周降级备用,最终平稳退役
核心经验总结:
渐进式迁移是大型中间件替换的唯一安全路径。四阶段策略(评估→并行→验证→切换)将风险逐级降低:评估阶段发现兼容性问题,并行阶段验证数据一致性,压测阶段证明性能达标,切换阶段确保平稳过渡。任何试图一步到位的方案都会将所有风险集中在单次切换上,失败代价极高。
分片策略映射工具化是效率关键。48个分片规则如果纯手工迁移,需要至少3周且容易出错。自动化映射工具将规则解析→类型映射→参数转换→YAML生成全流程自动化,耗时从3周缩短到2天。但自定义算法必须人工重写和验证——这部分占总工时的40%,无法自动化。
双写+一致性校验是并行阶段的基石。并行阶段的核心目标不是"两个中间件都可用",而是"数据绝对一致"。CRC校验每10分钟一次的全量对比,发现不一致时立即修复,这保证了并行阶段可以随时回退到MyCAT而无需担心数据缺失。
连接池调优不能凭经验,必须基于压测数据。HikariCP的默认参数对ShardingSphere Proxy场景不适用。4轮压测驱动的参数调优,最终找到适合我们业务负载的连接池配置。经验值可以作为起点,但终点必须由数据决定。
宽表预聚合是跨分片JOIN的性能杀手。分库分表中间件的跨分片JOIN天然有性能瓶颈,ShardingSphere比MyCAT快2倍多但仍不满足核心业务的延迟要求。宽表预聚合通过空间换时间,将实时JOIN转为预计算宽表查询,从根本上解决了跨分片JOIN的性能问题。
展望下一步,我们计划将ShardingSphere从Proxy模式切换到Sidecar模式(通过K8s DaemonSet部署),进一步降低中间件的网络开销和单点风险,同时为多集群联邦查询做好准备。
更多推荐
所有评论(0)