怎样用脚本批量迁移数据库?

wen 实用脚本 1

用Python脚本实现零停机迁移

怎样用脚本批量迁移数据库?

目录导读

  1. 为什么要用脚本批量迁移数据库?——痛点与场景分析
  2. 脚本迁移的核心技术栈选型对比
  3. 从零搭建批量迁移脚本(含完整代码示例)
  4. 关键问题问答:迁移中数据一致性、性能与回滚策略
  5. 生产环境迁移的最佳实践与避坑指南

为什么要用脚本批量迁移数据库?

在日常运维中,数据库迁移是高频操作,无论是从MySQL迁移到PostgreSQL,还是将SQL Server数据迁移到云原生数据库(如TiDB或GaussDB),手动操作不仅耗时,还可能因人为失误导致数据丢失,脚本批量迁移的核心优势在于:

  • 自动化:一次编写,多次执行,支持定时任务或CI/CD触发。
  • 可重复性:迁移逻辑固化在代码中,避免“记错表名”或“漏掉字段”。
  • 可审计:每一条数据的迁移状态都可以记录日志,便于事后追溯。
  • 零停机:通过增量同步策略,业务侧几乎无感知。

适用场景:数据库版本升级、跨云迁移、分库分表重构、同构/异构数据库复制。

脚本迁移的核心技术栈选型

方案 适用场景 语言/工具 优缺点
使用pandas+SQLAlchemy 中小型数据(<10GB) Python 灵活易调试,但内存占用高
PyMySQL+多线程分片 高并发大表 Python 速度快,需处理线程安全
pg_dump/pg_restore PostgreSQL => PostgreSQL Shell脚本 原生支持,但无法执行复杂转换
DataX(阿里巴巴) 异构数据库 Java+Python 稳定高效,但学习成本较高
推荐组合 通用场景 Python+PyMySQL+SQLAlchemy 兼顾性能与可维护性

核心原则:选择脚本语言时,优先考虑团队熟悉度,对于大数据量(>100GB),建议使用DataX或Apache Sqoop;对于中小规模,Python脚本足以满足需求。

从零搭建批量迁移脚本(核心代码示例)

假设我们需要将MySQL数据库中的orders表(1000万条数据)迁移到PostgreSQL,以下脚本实现了全量迁移+增量同步的核心逻辑:

# -*- coding: utf-8 -*-
import pymysql
import psycopg2
from datetime import datetime
import threading
from queue import Queue
# 数据库连接配置
mysql_config = {
    'host': '192.168.1.100',
    'port': 3306,
    'user': 'root',
    'password': 'password123',
    'database': 'source_db',
    'charset': 'utf8mb4'
}
pg_config = {
    'host': '192.168.1.200',
    'port': 5432,
    'user': 'admin',
    'password': 'pg_pass',
    'database': 'target_db'
}
# 批量迁移函数(分片读取)
def batch_migrate(table_name, chunk_size=10000):
    # 源库连接
    src_conn = pymysql.connect(**mysql_config)
    # 目标库连接
    dst_conn = psycopg2.connect(**pg_config)
    try:
        # 获取总行数
        with src_conn.cursor() as cursor:
            cursor.execute(f"SELECT COUNT(*) FROM {table_name}")
            total_rows = cursor.fetchone()[0]
        # 分批读取并写入
        offset = 0
        while offset < total_rows:
            with src_conn.cursor(pymysql.cursors.DictCursor) as cursor:
                cursor.execute(f"SELECT * FROM {table_name} LIMIT {chunk_size} OFFSET {offset}")
                rows = cursor.fetchall()
            # 批量插入pg
            with dst_conn.cursor() as pg_cursor:
                pg_cursor.executemany(
                    f"INSERT INTO {table_name} (id, order_date, amount) VALUES (%s, %s, %s) ON CONFLICT DO NOTHING",
                    [(row['id'], row['order_date'], row['amount']) for row in rows]
                )
            dst_conn.commit()
            offset += chunk_size
            print(f"已迁移 {min(offset, total_rows)} / {total_rows} 条")
    finally:
        src_conn.close()
        dst_conn.close()
# 多线程加速迁移(适用于大表)
def concurrent_migrate(table_name, thread_num=4):
    # 获取表的id范围
    src_conn = pymysql.connect(**mysql_config)
    with src_conn.cursor() as cursor:
        cursor.execute(f"SELECT MIN(id), MAX(id) FROM {table_name}")
        min_id, max_id = cursor.fetchone()
    src_conn.close()
    # 分片计算:每个线程负责的id段
    step = (max_id - min_id) // thread_num
    ranges = [(min_id + i*step, min_id + (i+1)*step - 1) for i in range(thread_num)]
    ranges[-1] = (ranges[-1][0], max_id)
    def worker(id_range):
        # 每个线程独立处理自己的id段
        src_conn = pymysql.connect(**mysql_config)
        dst_conn = psycopg2.connect(**pg_config)
        try:
            with src_conn.cursor(pymysql.cursors.DictCursor) as cursor:
                cursor.execute(f"SELECT * FROM {table_name} WHERE id BETWEEN {id_range[0]} AND {id_range[1]}")
                rows = cursor.fetchall()
            with dst_conn.cursor() as pg_cursor:
                pg_cursor.executemany(
                    f"INSERT INTO {table_name} VALUES (%(id)s, %(order_date)s, %(amount)s)",
                    rows
                )
            dst_conn.commit()
        finally:
            src_conn.close()
            dst_conn.close()
    threads = [threading.Thread(target=worker, args=(r,)) for r in ranges]
    for t in threads: t.start()
    for t in threads: t.join()
    print("多线程迁移完成")

说明

  • 分片读取:通过LIMIT/OFFSET避免内存溢出。
  • 冲突处理:使用ON CONFLICT DO NOTHING保障幂等性。
  • 多线程分片:根据主键范围拆分任务,提升迁移速度3-5倍。

关键问题问答

Q1:迁移过程中如何保证数据一致性? A:采用“先全量后增量”策略,全量迁移开始时记录当前时间戳T1,全量完成后开启增量同步(如MySQL的binlog解析或数据库的CDC机制),将T1之后产生的变更同步到目标库,最终在业务低峰期进行短时停机验证,确认两边数据完全一致。

Q2:迁移脚本性能瓶颈在哪,如何优化? A:常见瓶颈包括:

  • 源库IO:大表全量扫描会导致源库磁盘IO飙升,建议在业务低峰期执行。
  • 网络延迟:批量插入时使用executemany(Python)或批量SQL语句,减少网络往返。
  • 目标索引:建议先关闭目标表的索引和约束,迁移完成后重建(可提升50%以上写入速度)。

Q3:迁移失败后如何回滚? A:至少准备两条路径:

  • 脚本级回滚:记录每个分片的迁移状态(如使用Redis或MySQL状态表),失败时仅重试失败分片。
  • 数据库级回滚:在目标库创建“前置快照”(云数据库支持闪回),或保留完整的迁移日志(如CSV文件)以便手动重建。

生产环境迁移的最佳实践

  1. 数据校验先行:迁移前验证源库和目标库的字符集、字段类型是否兼容,MySQL的DATETIME映射到PostgreSQL时需注意时区转换。
  2. 安全第一:脚本中明文密码需替换为环境变量(如os.getenv('DB_PASS'))。
  3. 限流保护:在迁移循环中加入time.sleep(0.1),避免对源库造成过大压力。
  4. 日志全量记录:每条迁移记录输出表名、主键ID、时间戳到文件,便于异常分析。
  5. 灰度验证:先迁移5%的数据(如2023年的订单),验证业务无异常后再执行全量迁移。

避坑提醒

  • 不要信任“一键迁移”工具,尤其是跨数据库异构迁移,字段类型转换常出现精度丢失(如DECIMAL(10,2)FLOAT导致四舍五入差异)。
  • 迁移完成后必须删除脚本中的临时文件(如保存密码的配置文件),防止安全漏洞。

抱歉,评论功能暂时关闭!