
数据迁移的十大致命错误半年事故复盘的避坑指南数据迁移是数据库工作中风险最高的操作之一。过去半年团队经历和见证了多起数据迁移事故有的导致了数小时的服务中断有的造成了数据的静默丢失。本文复盘这些事故的共性根因提炼出十大避坑原则。一、当2TB数据迁移让业务中断8小时一次典型事故的全景回顾今年1月一个核心业务库需要从自建IDC迁移到云端。迁移方案看起来无懈可击使用DTS全量增量同步灰度切流预计停机时间30分钟。实际执行中出现了连环问题全量同步阶段源库一个2TB的大表包含了一个隐藏的大字段DTS默认的读取超时设置不足导致全量同步中断了3次增量同步阶段源库的binlog格式是MIXED一条UPDATE语句被记录为STATEMENT格式目标库执行时因为时区设置不同导致数据更新到了错误的时间范围切流后应用层连接池未正确切换到新库部分请求仍然打到老库产生了双写导致数据冲突最终业务中断8小时修复数据不一致又花了两天时间。复盘发现每一个问题都可以通过更完善的事前检查来避免。二、数据迁移的事故模式与根因映射三、数据迁移安全工具箱#!/usr/bin/env python3 Data Migration Safety Check Toolkit import hashlib import pymysql import time import json from typing import Dict, List, Tuple, Optional from dataclasses import dataclass from concurrent.futures import ThreadPoolExecutor dataclass class MigrationCheckResult: check_name: str passed: bool details: str severity: str # BLOCKER, WARNING, INFO class MigrationSafetyChecker: def __init__(self, source_config: dict, target_config: dict): self.source_config source_config self.target_config target_config self.results: List[MigrationCheckResult] [] def _connect(self, config: dict) - Optional[pymysql.Connection]: try: return pymysql.connect(**config) except pymysql.Error as e: print(f[ERROR] Connection failed: {e}) return None def check_version_compatibility(self): 检查1: 版本兼容性 source self._connect(self.source_config) target self._connect(self.target_config) try: src_ver None tgt_ver None if source: with source.cursor() as cur: cur.execute(SELECT VERSION()) src_ver cur.fetchone()[0] if target: with target.cursor() as cur: cur.execute(SELECT VERSION()) tgt_ver cur.fetchone()[0] major_src int(src_ver.split(.)[0]) if src_ver else 0 major_tgt int(tgt_ver.split(.)[0]) if tgt_ver else 0 if major_src major_tgt: self.results.append(MigrationCheckResult( 版本兼容性, True, f源库{src_ver} - 目标库{tgt_ver}升级迁移, INFO )) elif major_src major_tgt: self.results.append(MigrationCheckResult( 版本兼容性, True, f同版本迁移 {src_ver}, INFO )) else: self.results.append(MigrationCheckResult( 版本兼容性, False, f降级迁移不支持 {src_ver} - {tgt_ver}, BLOCKER )) except Exception as e: self.results.append(MigrationCheckResult( 版本兼容性, False, str(e), BLOCKER )) finally: if source: source.close() if target: target.close() def check_character_set(self): 检查2: 字符集一致性 source self._connect(self.source_config) target self._connect(self.target_config) try: charsets {} for label, conn in [(source, source), (target, target)]: if conn: with conn.cursor() as cur: cur.execute(SHOW VARIABLES LIKE character_set_%) charsets[label] {r[0]: r[1] for r in cur.fetchall()} for key in [character_set_server, character_set_database]: src_val charsets.get(source, {}).get(key) tgt_val charsets.get(target, {}).get(key) if src_val and tgt_val and src_val ! tgt_val: self.results.append(MigrationCheckResult( 字符集一致性, False, f{key}: 源{src_val}, 目标{tgt_val}, WARNING )) except Exception as e: self.results.append(MigrationCheckResult( 字符集一致性, False, str(e), WARNING )) finally: if source: source.close() if target: target.close() def check_binlog_format(self): 检查3: Binlog格式DTS增量同步必需ROW格式 source self._connect(self.source_config) try: if source: with source.cursor() as cur: cur.execute(SHOW VARIABLES LIKE binlog_format) result cur.fetchone() if result: fmt result[1] passed fmt ROW self.results.append(MigrationCheckResult( Binlog格式, passed, f当前格式: {fmt} ( (满足DTS要求) if passed else (DTS增量同步需要ROW格式!)), BLOCKER if not passed else INFO )) except Exception as e: self.results.append(MigrationCheckResult( Binlog格式, False, str(e), BLOCKER )) finally: if source: source.close() def check_large_tables(self, size_threshold_gb: float 100): 检查4: 大表检测 source self._connect(self.source_config) try: if source: with source.cursor() as cur: cur.execute( SELECT table_schema, table_name, ROUND((data_length index_length) / 1024 / 1024 / 1024, 2) as size_gb, table_rows FROM information_schema.tables WHERE table_schema NOT IN (mysql, sys, performance_schema, information_schema) AND (data_length index_length) / 1024 / 1024 / 1024 %s ORDER BY size_gb DESC , (size_threshold_gb,)) for row in cur.fetchall(): self.results.append(MigrationCheckResult( 大表检测, True, f{row[0]}.{row[1]}: {row[2]}GB, {row[3]} 行 - 需要分批迁移, WARNING )) except Exception as e: self.results.append(MigrationCheckResult( 大表检测, False, str(e), WARNING )) finally: if source: source.close() def check_triggers_and_procedures(self): 检查5: 触发器/存储过程/函数/事件 source self._connect(self.source_config) try: if source: with source.cursor() as cur: # 触发器 cur.execute( SELECT COUNT(*) FROM information_schema.triggers WHERE trigger_schema NOT IN (mysql, sys) ) trigger_count cur.fetchone()[0] # 存储过程 cur.execute( SELECT COUNT(*) FROM information_schema.routines WHERE routine_schema NOT IN (mysql, sys) ) routine_count cur.fetchone()[0] if trigger_count 0: self.results.append(MigrationCheckResult( 触发器/存储过程, True, f检测到{trigger_count}个触发器, {routine_count}个存储过程 - DTS不会自动迁移,需手动处理, WARNING )) except Exception as e: self.results.append(MigrationCheckResult( 触发器/存储过程, False, str(e), WARNING )) finally: if source: source.close() def check_no_pk_tables(self): 检查6: 无主键表检测DTS增量同步必需主键 source self._connect(self.source_config) try: if source: with source.cursor() as cur: cur.execute( SELECT t.table_schema, t.table_name FROM information_schema.tables t LEFT JOIN information_schema.table_constraints c ON t.table_schema c.table_schema AND t.table_name c.table_name AND c.constraint_type PRIMARY KEY WHERE t.table_schema NOT IN (mysql, sys, performance_schema, information_schema) AND t.table_type BASE TABLE AND c.constraint_name IS NULL ) no_pk_tables cur.fetchall() for row in no_pk_tables: self.results.append(MigrationCheckResult( 无主键表, False, f{row[0]}.{row[1]} 缺少主键 - DTS增量同步无法处理, BLOCKER )) except Exception as e: self.results.append(MigrationCheckResult( 无主键表, False, str(e), BLOCKER )) finally: if source: source.close() def run_all_checks(self) - List[MigrationCheckResult]: 运行所有检查 checks [ self.check_version_compatibility, self.check_character_set, self.check_binlog_format, self.check_large_tables, self.check_triggers_and_procedures, self.check_no_pk_tables, ] with ThreadPoolExecutor(max_workers4) as executor: futures [executor.submit(check) for check in checks] for future in futures: try: future.result(timeout30) except Exception as e: print(f[ERROR] Check timeout/error: {e}) return self.results def print_report(self): 输出检查报告 blockers [r for r in self.results if r.severity BLOCKER] warnings [r for r in self.results if r.severity WARNING] print(\n * 70) print( 数据迁移安全检查报告) print( * 70) print(f BLOCKER: {len(blockers)} | WARNING: {len(warnings)}) print(- * 70) for r in self.results: symbol PASS if r.passed else FAIL print(f [{r.severity}] [{symbol}] {r.check_name}) print(f {r.details}) if blockers: print(\n[BLOCKER] 存在阻塞性问题迁移不能继续执行!) else: print(\n[OK] 阻塞性检查通过可以继续迁移流程) if __name__ __main__: checker MigrationSafetyChecker( source_config{ host: source-db.example.com, user: migration_check, password: your_password, charset: utf8mb4 }, target_config{ host: target-db.example.com, user: migration_check, password: your_password, charset: utf8mb4 } ) checker.run_all_checks() checker.print_report()四、十大致命错误的速查清单序号致命错误后果预防措施1不检查binlog格式增量同步失败确认binlog_formatROW2忽略无主键表增量数据无法同步迁移前给所有表加主键3字符集不统一数据乱码/静默丢失确保character_set一致4大表一次性全量同步源库CPU打满按时间/ID范围分批5未做数据校验数据不一致不自知pt-table-checksum 业务校验6遗漏存储过程和触发器业务逻辑缺失导出后手动在目标库创建7忽略连接池切换新旧库双写制定详细的切流SOP8未准备回滚方案故障后无路可退保留源库只读回切DNS9高峰期执行迁移用户感知故障选择业务低峰提前公告10忽略大字段影响同步超时中断提前识别BLOB/TEXT大字段五、总结数据迁移的本质不是技术操作而是风险管理。每一条经验教训背后都是一次刻骨铭心的故障。如果只能记住一条原则那就是永远不要相信应该没问题——在正式执行前对所有假设做最小化验证。下半年将在自动化迁移检查工具链上加大投入目标是让每一项检查都自动化执行将人为疏忽的概率降到最低。