怎样实现分批迁移海量数据库数据

wen 实用脚本 32

本文目录导读:

怎样实现分批迁移海量数据库数据

  1. 核心原则
  2. 方案一:基于主键或唯一索引的范围分片(最高效)
  3. 方案二:基于游标(Cursor)的分页(适用于非自增ID)
  4. 方案三:使用数据库的COPY命令(最快速)
  5. 方案四:并行迁移(提高吞吐量)
  6. 方案五:使用ETL工具(省心但需运维)
  7. 关键注意事项
  8. 总结选择建议

实现海量数据库数据的分批迁移,核心思路是避免一次性加载所有数据到内存,而是通过分页、游标或分区的方式,将数据切成多个小块(Chunk),逐块处理。

下面介绍几种主流且高效的实现方案,包括原理、代码示例(以Python + MySQL为例)以及注意事项。

核心原则

  • *避免`SELECT OFFSET LIMIT大偏移量**:传统的LIMIT 1000000, 1000`方式随着偏移量增大,性能会急剧下降(数据库需要扫描并丢弃前面的大量行)。
  • 利用索引:分批的关键是高效的定位“下一条该取什么”。
  • 事务与锁:迁移过程中需考虑数据一致性,可能需要结合SELECT ... FOR UPDATE或使用数据库的快读(如MySQL的Repeatable Read级别)。

基于主键或唯一索引的范围分片(最高效)

这是处理海量数据(千万级、亿级)最推荐的方式,利用表的主键(通常是自增ID或连续有序的ID)。

原理: 估算一条记录的大小,设定一个“页大小”(如每批5000条),通过记录当前处理的最大ID来获取下一批数据。

Python 示例(伪代码)

import pymysql
from typing import Any, Generator
BATCH_SIZE = 5000  # 每批处理5000条,可根据单条大小调整
def fetch_data_in_batches(source_cursor, target_cursor, start_id: int, end_id: int):
    """从MySQL分批读取并写入目标库"""
    current_start = start_id
    while current_start < end_id:
        current_end = current_start + BATCH_SIZE
        # 关键SQL:使用 WHERE id BETWEEN ... AND ...
        sql = f"""
            SELECT * FROM source_table 
            WHERE id >= %s AND id < %s 
            ORDER BY id ASC
        """
        source_cursor.execute(sql, (current_start, current_end))
        # 注意:这里不要一次性fetchall(),除非BATCH_SIZE很小且内存够
        # 更好的处理是:逐行处理并批量写入目标
        rows = source_cursor.fetchmany(size=BATCH_SIZE)
        if not rows:
            break
        # 批量写入目标库(使用executemany或批量insert)
        target_cursor.executemany(
            "INSERT INTO target_table (col1, col2, ...) VALUES (%s, %s, ...)", 
            rows
        )
        # 提交当前批次
        target_connection.commit()
        print(f"已迁移 ID 范围: {current_start} - {current_end-1},共 {len(rows)} 条")
        # 更新起始ID为当前批次的结束ID
        current_start = current_end
    # 处理最后一批不足BATCH_SIZE的情况
    if rows:
        pass # 已处理
# 使用示例
source_conn = pymysql.connect(...)
target_conn = pymysql.connect(...)
with source_conn.cursor() as src_cursor, target_conn.cursor() as tgt_cursor:
    # 获取表的ID范围(需要提前获取)
    src_cursor.execute("SELECT MIN(id), MAX(id) FROM source_table")
    min_id, max_id = src_cursor.fetchone()
    fetch_data_in_batches(src_cursor, tgt_cursor, min_id, max_id)

优点

  • 性能稳定,不会随着数据量增加而变慢。
  • 容易断点续传(记录最后处理的ID即可)。

缺点

  • 必须依赖有序且连续的主键,如果ID有大量空洞(删除导致),仍然不影响,因为BETWEEN会跳过空洞,但需要确保即使空洞很大,下一批也能正确衔接(因为是根据id < 来取的,空洞后的ID更大,会落入下一批)。
  • 如果主键不是有序的(如UUID),此方案效果差。

基于游标(Cursor)的分页(适用于非自增ID)

当主键不是简单的自增ID时,使用“键值游标法”(Keyset Pagination 或 Seek Method)。

原理:不依赖OFFSET,而是利用WHERE条件直接定位到上一批的最后一条记录。

-- 第一页:取前 1000 条
SELECT * FROM source_table ORDER BY create_time, id LIMIT 1000;
-- 第二页:记住上一页最后一条记录的 create_time 和 id
SELECT * FROM source_table 
WHERE (create_time > '上一页最后时间') 
   OR (create_time = '上一页最后时间' AND id > '上一页最后ID')
ORDER BY create_time, id 
LIMIT 1000;

Python 示例

def migrate_by_keyset(cursor, last_values: tuple):
    """
    last_values: 上一批最后一条记录的排序字段值及主键值,假设排序字段为 (create_time, id)
    """
    if last_values is None:
        # 第一批:直接取
        sql = "SELECT * FROM source_table ORDER BY create_time, id LIMIT %s"
        cursor.execute(sql, (BATCH_SIZE,))
    else:
        # 后续批次:使用游标定位
        last_create_time, last_id = last_values
        sql = """
            SELECT * FROM source_table 
            WHERE (create_time > %s OR (create_time = %s AND id > %s))
            ORDER BY create_time, id 
            LIMIT %s
        """
        cursor.execute(sql, (last_create_time, last_create_time, last_id, BATCH_SIZE))
    rows = cursor.fetchall()
    if rows:
        # 记录最后一条记录用于下一批
        last_row = rows[-1]
        next_cursor = (last_row['create_time'], last_row['id'])
    else:
        next_cursor = None
    return rows, next_cursor

优点

  • 对任何有唯一索引或组合索引的表都有效。
  • 无论数据多大,每次查询都只扫描LIMIT大小的数据,性能极佳。

缺点

  • 需要确保排序字段有索引,且排序字段组合保证唯一性(否则会漏/重复)。
  • 实现稍微复杂一点,需要维护游标状态。

使用数据库的COPY命令(最快速)

如果目标也是同一个数据库类型(如MySQL到MySQL,PostgreSQL到PostgreSQL),利用原生工具是最快的。

  • MySQLSELECT ... INTO OUTFILE + LOAD DATA INFILE
    • 导出:SELECT * FROM table WHERE id BETWEEN ... INTO OUTFILE '/tmp/part_1.csv' FIELDS TERMINATED BY ',';
    • 导入:LOAD DATA INFILE '/tmp/part_1.csv' INTO TABLE target_table FIELDS TERMINATED BY ',';
  • PostgreSQLCOPY 命令。
  • 云数据库:使用DTS(数据传输服务)等托管工具,支持并行迁移。

并行迁移(提高吞吐量)

对于超大表(如10亿条),单线程太慢,可以将ID范围拆成多个独立的区间,开多个进程/线程并行处理。

思路

  1. 获取MIN(id)MAX(id)
  2. 计算总数据量,按N个并发任务,将ID范围分成N等份(假设ID相对均匀)。
  3. 每个线程独立执行方案一或方案二,处理自己的ID段。

注意并发控制

  • 目标端:批量写入时,使用INSERT ... ON DUPLICATE KEY UPDATE(防止重复)或者分批提交。
  • 源端:读操作天然支持并发,但要注意数据库连接数限制。

使用ETL工具(省心但需运维)

如果不想写大量代码,可以使用成熟工具:

  • DataX (阿里):支持多种数据源,自带分片算法。
  • Kettle (Pentaho):图形化界面,支持分页和记录集限制。
  • Apache Flink / Spark:处理实时或历史数据迁移,适合复杂转换。

关键注意事项

  1. 内存管理:绝对不要fetchall()一次拉取所有数据,使用fetchmany(size)或逐行处理。
  2. 事务控制
    • 源端:如果迁移过程中源表还在写,需要考虑使用SELECT ... FOR UPDATE(影响性能)或使用数据库的Binlog/Flink CDC进行实时同步。
    • 目标端:建议每批提交一次(batch_commit),而不是每行或最后一次性提交,每批提交后,如果进程崩溃,最多丢失当前这一批数据(可重跑)。
  3. 断点续传
    • 方案一:记录last_processed_id
    • 方案二:记录上次游标值。
    • 建议将进度写入到一个持久化位置(如数据库表、Redis、文件),以便失败重启时跳过已完成的批次。
  4. 性能监控
    • 开启慢查询日志。
    • 监控源库的IOPSCPU,不要在业务高峰期运行。
    • 分批写入目标时,适当调整BATCH_SIZE(通常500-5000条/批,或每批数据量不超过1MB),避免目标库日志暴增或锁表。
  5. 数据类型
    • 注意大字段(TEXTBLOB),如果包含大字段,批次应该更小(如200条),否则可能内存溢出或网络超时。
    • 注意时区、字符集转换。
  6. 网络:如果源和目标在不同网络/机房,考虑使用压缩传输。
  7. 目标库约束:迁移前禁用目标表的外键检查唯一性检查(如SET FOREIGN_KEY_CHECKS=0; SET UNIQUE_CHECKS=0;),完成后再启用,可以大幅提高插入速度。

总结选择建议

  • 表有自增ID,数据量大,网络好方案一(范围分片) + 并行迁移
  • 表没有自增ID,但有其他有序字段方案二(游标分页)
  • 纯数据库到数据库,且在同一网络 → 考虑 COPY / LOAD DATA 方案。
  • 需要转换数据结构(如MySQL到MongoDB) → 使用 DataXApache Spark
  • 需要不停机迁移 → 结合 Binlog (CDC) + 数据全量迁移 + 增量追同步(这是更高级的方案,但正确性最高)。

务必先在测试环境验证:用生产数据的1/100测试一下,确认数据完整性、迁移速度、内存消耗,然后再运行。

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