昨天深夜两点,我的手机响了。不是闹钟,是公司监控报警群的消息——“核心交易表数据异常,主从不一致率0.3%”。0.3%听起来不多?但那个表里躺着的是几千万条用户订单记录,每一笔都连着真金白银。
我爬起来打开笔记本,连上VPN,开始了一场持续了整整六个小时的排查。今天就把这个过程完整记录下来,希望能帮到正在经历同样困境的你。
一、问题是怎么被发现的
一切始于用户投诉。一位VIP客户反映他在手机上看到的订单金额和客服那边查到的不一样。这可不是小事,尤其是在金融相关的应用里,数据不一致意味着信任崩塌。
我们第一反应是查应用层日志,但应用日志显示两个接口返回的数据是一致的——这反而让问题变得更诡异。既然应用层没bug,那问题肯定出在数据库层面。
1.1 初步排查:确认不一致范围
我写了个简单的对比脚本,从主库和从库各取了一批数据做diff:
-- 在主库执行
SELECT id, amount, updated_at
FROM orders
WHERE updated_at > '2024-01-15 00:00:00'
ORDER BY id
INTO OUTFILE '/tmp/master_orders.csv'
FIELDS TERMINATED BY ','
ENCLOSED BY '"'
LINES TERMINATED BY '\n';
-- 在从库执行同样的查询
SELECT id, amount, updated_at
FROM orders
WHERE updated_at > '2024-01-15 00:00:00'
ORDER BY id
INTO OUTFILE '/tmp/slave_orders.csv'
FIELDS TERMINATED BY ','
ENCLOSED BY '"'
LINES TERMINATED BY '\n';
拿到两个CSV文件后,我用Python简单对比了一下:
import csv
def load_csv(filepath):
data = {}
with open(filepath, 'r', encoding='utf-8') as f:
reader = csv.DictReader(f)
for row in reader:
data[row['id']] = {
'amount': row['amount'],
'updated_at': row['updated_at']
}
return data
master_data = load_csv('/tmp/master_orders.csv')
slave_data = load_csv('/tmp/slave_orders.csv')
mismatches = []
for id, master_record in master_data.items():
if id not in slave_data:
mismatches.append({
'id': id,
'type': 'slave_missing',
'master': master_record
})
elif master_record != slave_data[id]:
mismatches.append({
'id': id,
'type': 'data_diff',
'master': master_record,
'slave': slave_data[id]
})
print(f"发现 {len(mismatches)} 条不一致记录")
for m in mismatches[:10]:
print(m)
结果显示有将近两千条记录不一致,主要集中在最近两个小时。这很关键——说明问题是近期发生的,不是历史遗留。
1.2 检查主从复制状态
接下来我登录到MySQL,检查了主从复制的状态:
-- 在主库上执行
SHOW MASTER STATUS\G
-- 在从库上执行
SHOW SLAVE STATUS\G
关键信息如下:
Slave_IO_Running: Yes
Slave_SQL_Running: Yes
Seconds_Behind_Master: 45
Relay_Log_Space: 254789632
Last_Error:
Seconds_Behind_Master: 45 说明从库已经落后主库45秒了。这个数字看着不大,但在高并发写入的场景下,45秒意味着可能有成千上万条更新还没同步过去。
更让我警惕的是,Relay_Log_Space 这个值比平时的基线高了将近30%,说明中继日志堆积得比较厉害。
二、深入分析:延迟从哪来的
确认了有延迟,下一步就是搞清楚延迟产生的原因。 MySQL的主从复制流程其实很清晰:主库写Binlog → 从库IO线程拉取Binlog到本地中继日志 → 从库SQL线程重放中继日志。任何一个环节卡住都会导致延迟。
2.1 排除网络问题
我首先排除了网络因素。主从服务器之间是专线连接,ping值稳定在1ms以内,带宽也很充裕。而且如果是网络问题,IO线程应该会报错或者断开,但状态显示IO线程是正常运行的。
2.2 检查SQL线程是否成为瓶颈
这是我最怀疑的点。我查看了从库的实时负载:
top -bn1 | head -20
发现CPU使用率并不高,内存也充足。这说明不是硬件资源不足导致的问题。
那我转而检查从库是否在执行大事务或者锁表:
-- 检查当前正在执行的查询
SHOW PROCESSLIST;
-- 检查锁等待情况
SELECT * FROM information_schema.INNODB_TRX;
SELECT * FROM information_schema.INNODB_LOCK_WAITS;
PROCESSLIST显示从库上有一个正在执行的ALTER TABLE操作!这个操作占据了大量的CPU和IO资源,导致SQL线程无法及时重放Binlog。
Id: 1234
User: root
Host: localhost
db: production
Command: Query
Time: 8432
State: Sending data
Info: ALTER TABLE orders ADD INDEX idx_status_created (status, created_at)
8432秒!这个ALTER TABLE已经执行了将近两个半小时!这就是延迟的罪魁祸首。
2.3 验证推断
为了确认我的推断,我查看了这个ALTER TABLE操作的起始时间:
-- 查看近期执行过的DDL
SELECT * FROM mysql.general_log
WHERE argument LIKE '%ALTER TABLE orders%'
ORDER BY event_time DESC
LIMIT 5;
确认这条DDL是在问题发生前大约一小时开始执行的。时间线完全吻合。
三、紧急处理:如何快速恢复
找到问题根源后,接下来的任务是如何快速恢复主从同步。我采取了以下几个步骤:
3.1 暂停大事务(如果可能)
首先尝试暂停那个ALTER TABLE操作。但在MySQL中,ALTER TABLE是不可中断的,只能等待它完成。这是一个痛苦的等待过程。
3.2 评估数据一致性风险
在等待的同时,我必须评估数据不一致的风险。我写了一个更详细的对比脚本,检查不一致记录的业务影响:
import pymysql
import json
def compare_db_consistency():
# 连接主库和从库
master_conn = pymysql.connect(
host='master-db.internal',
user='readonly',
password='***',
database='production',
charset='utf8mb4'
)
slave_conn = pymysql.connect(
host='slave-db.internal',
user='readonly',
password='***',
database='production',
charset='utf8mb4'
)
master_cur = master_conn.cursor()
slave_cur = slave_conn.cursor()
# 获取最近1小时的变更
cutoff_time = '2024-01-15 02:00:00'
# 查询主库
master_cur.execute(f"""
SELECT id, order_no, amount, status, updated_at
FROM orders
WHERE updated_at > '{cutoff_time}'
""")
master_records = {row[0]: row for row in master_cur.fetchall()}
# 查询从库
slave_cur.execute(f"""
SELECT id, order_no, amount, status, updated_at
FROM orders
WHERE updated_at > '{cutoff_time}'
""")
slave_records = {row[0]: row for row in slave_cur.fetchall()}
# 对比
inconsistencies = []
for id, master_row in master_records.items():
if id not in slave_records:
inconsistencies.append({
'id': id,
'type': 'missing_in_slave',
'data': master_row
})
else:
slave_row = slave_records[id]
# 比较关键字段
if master_row[2] != slave_row[2] or master_row[3] != slave_row[3]: # amount, status
inconsistencies.append({
'id': id,
'type': 'field_diff',
'master': master_row,
'slave': slave_row
})
# 统计分析
missing_count = sum(1 for i in inconsistencies if i['type'] == 'missing_in_slave')
diff_count = sum(1 for i in inconsistencies if i['type'] == 'field_diff')
print(f"总不一致记录: {len(inconsistencies)}")
print(f"从库缺失: {missing_count}")
print(f"字段差异: {diff_count}")
# 检查金额差异的总额
total_amount_diff = 0
for i in inconsistencies:
if i['type'] == 'field_diff':
diff = abs(float(i['master'][2]) - float(i['slave'][2]))
total_amount_diff += diff
print(f"金额差异总额: ¥{total_amount_diff:.2f}")
master_conn.close()
slave_conn.close()
return inconsistencies
inconsistencies = compare_db_consistency()
运行结果显示,金额差异总额超过了50万元。这是一个必须立即处理的严重问题。
3.3 决定处理策略
我有两个选择:
- 等待ALTER TABLE完成:可能需要再等几个小时,期间数据不一致风险持续存在。
- 强制重置主从复制:丢弃从库当前数据,重新同步。但这会导致从库短暂不可用,且可能丢失部分数据。
经过和业务方沟通,我们决定采取折中方案:在从库上创建一个新表,将主库的最新数据导入,然后切换读取流量到新表。这样既保证了一致性,又避免了等待ALTER TABLE完成。
3.4 实施数据同步
# 在主库上导出最新数据
mysqldump -h master-db.internal -u admin -p production orders \
--where="updated_at > '2024-01-15 00:00:00'" \
--single-transaction \
--quick \
--lock-tables=false \
> orders_incremental.sql
# 在从库上创建新表
mysql -h slave-db.internal -u admin -p production < create_new_table.sql
# 导入数据
mysql -h slave-db.internal -u admin -p production production < orders_incremental.sql
其中create_new_table.sql的内容:
CREATE TABLE orders_latest AS SELECT * FROM orders LIMIT 0;
ALTER TABLE orders_latest ADD INDEX idx_status_created (status, created_at);
-- 复制原表的其他索引和约束
导入完成后,我验证了新表的数据一致性:
-- 在主库和从库上分别执行
SELECT COUNT(*) FROM orders_latest;
SELECT SUM(amount) FROM orders_latest WHERE updated_at > '2024-01-15 00:00:00';
两边结果完全一致,数据同步成功。
3.5 切换流量并监控
接下来是切换读取流量。我们在应用配置中增加了数据源切换的逻辑:
// Spring配置
@Bean
public DataSource routingDataSource(@Qualifier("masterDataSource") DataSource master,
@Qualifier("slaveDataSource") DataSource slave) {
RoutingDataSource routingDataSource = new RoutingDataSource();
Map<Object, Object> targetDataSources = new HashMap<>();
targetDataSources.put(DatabaseType.MASTER, master);
targetDataSources.put(DatabaseType.SLAVE, slave);
routingDataSource.setTargetDataSources(targetDataSources);
routingDataSource.setDefaultTargetDataSource(master);
return routingDataSource;
}
// 动态切换数据源
public void switchToLatestSlave() {
// 创建一个新的数据源指向最新表
DataSource latestDataSource = createDataSourceForLatestTable();
routingDataSource.setDataSources(
Map.of(DatabaseType.SLAVE, latestDataSource)
);
routingDataSource.afterPropertiesSet();
}
切换完成后,我密切监控了半小时,确认没有新的不一致产生,主从复制状态逐渐恢复正常。
四、长期方案:双写校验机制
这次事故让我意识到,仅仅依靠MySQL自带的主从复制是不够的。我们需要一个更主动的校验机制,能够在问题发生时就发现并告警,而不是等到用户投诉。
4.1 设计思路
我设计了一个双写校验方案,核心思路是:
- 定期抽样对比:每隔一段时间,随机抽取一批数据,在主库和从库上进行对比。
- 关键业务实时校验:对于核心表,每次写操作后都进行校验。
- 自动修复机制:发现不一致时,自动从主库同步数据到从库。
4.2 实现细节
4.2.1 抽样对比服务
”`python import pymysql import random import logging from datetime import datetime, timedelta import hashlib
logging.basicConfig(level=logging.INFO) logger = logging.getLogger(name)
class DataConsistencyChecker:
def __init__(self, master_config, slave_config):
self.master_conn = pymysql.connect(**master_config)
self.slave_conn = pymysql.connect(**slave_config)
self.master_cur = self.master_conn.cursor()
self.slave_cur = self.slave_conn.cursor()
def hash_record(self, record):
"""生成记录的唯一哈希值"""
record_str = '|'.join(str(v) for v in record if v is not None)
return hashlib.md5(record_str.encode()).hexdigest()
def check_table(self, table_name, sample_size=1000, where_clause=None):
"""对指定表进行抽样对比"""
# 构建查询
base_query = f"SELECT * FROM {table_name}"
if where_clause:
base_query += f" WHERE {where_clause}"
# 从主库随机抽样
sample_query = f"{base_query} ORDER BY RAND() LIMIT {sample_size}"
self.master_cur.execute(sample_query)
master_samples = self.master_cur.fetchall()
# 获取主键列名
self.master_cur.execute(f"DESC {table_name}")
columns = [col[0] for col in self.master_cur.fetchall()]
primary_keys = self._get_primary_keys(table_name)
if not primary_keys:
logger.warning(f"表 {table_name} 没有主键,跳过一致性检查")
return []
# 构建主键到记录的映射
master_records = {}
for row in master_samples:
key_values = tuple(row[columns.index(pk)] for pk in primary_keys)
master_records[key_values] = row
# 从从库查询相同记录
inconsistencies = []
for key_values in master_records.keys():
placeholders = ','.join(['%s'] * len(primary_keys))
self.slave_cur.execute(
f"SELECT * FROM {table_name} WHERE {self._build_where_clause(primary_keys)}",
key_values
)
slave_record = self.slave_cur.fetchone()
if slave_record is None:
inconsistencies.append({
'type': 'missing_in_slave',
'key': key_values,
'master_record': master_records[key_values]
})
else:
# 对比记录内容(排除时间戳等可能变化的字段)
if self._records_differ(master_records[key_values], slave_record, columns):
inconsistencies.append({
'type': 'field_diff',
'key': key_values,
'master_record': master_records[key_values],
'slave_record': slave_record
})
return inconsistencies
def _get_primary_keys(self, table_name):
"""获取表的主键列"""
self.master_cur.execute(f"""
SELECT COLUMN_NAME
FROM INFORMATION_SCHEMA.KEY_COLUMN_USAGE
WHERE TABLE_SCHEMA = DATABASE()
AND TABLE_NAME = %s
AND CONSTRAINT_NAME = 'PRIMARY'
""", (table_name,))
return [row[0] for row in self.master_cur.fetchall()]
def _build_where_clause(self, primary_keys):
"""构建WHERE子句"""
return ' AND '.join([f"`{pk}` = %s" for pk in primary_keys])
def _records_differ(self, master_record, slave_record, columns):
"""比较两条记录是否不同(排除时间戳字段)"""
skip_columns = {'updated_at', 'created_at', 'update_time', 'create_time'}
for i, col in enumerate(columns):
if col in skip_columns:
continue
if master_record[i] != slave_record[i]:
return True
return False
def run_daily_check(self):
"""执行每日全量检查"""
logger.info("开始执行每日数据一致性检查...")
# 定义需要检查的核心表
tables_to_check = [
('orders', 'updated_at > DATE_SUB(NOW(), INTERVAL 24 HOUR)'),
('users', 'updated_at > DATE_SUB(NOW(), INTERVAL 24 HOUR)'),
('payments', 'updated_at > DATE_SUB(NOW(), INTERVAL 24 HOUR)'),
]
all_inconsistencies = []
for table_name, where_clause in tables_to_check:
logger.info(f"检查表 {table_name}...")
inconsistencies = self.check_table(table_name, sample_size=500, where_clause=where_clause)
if inconsistencies:
logger.warning(f"表 {table_name} 发现 {len(inconsistencies)} 条不一致记录")
all_inconsistencies.extend(inconsistencies)
# 发送告警
self._send_alert(table_name, inconsistencies)
else:
logger.info(f"表 {table_name} 数据一致")
# 生成报告
if all_inconsistencies:
self._generate_report(all_inconsist
