本文目录导读:

- 核心原则
- 方案一:基于主键或唯一索引的范围分片(最高效)
- 方案二:基于游标(Cursor)的分页(适用于非自增ID)
- 方案三:使用数据库的
COPY命令(最快速) - 方案四:并行迁移(提高吞吐量)
- 方案五:使用ETL工具(省心但需运维)
- 关键注意事项
- 总结选择建议
实现海量数据库数据的分批迁移,核心思路是避免一次性加载所有数据到内存,而是通过分页、游标或分区的方式,将数据切成多个小块(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),利用原生工具是最快的。
- MySQL:
SELECT ... 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 ',';
- 导出:
- PostgreSQL:
COPY命令。 - 云数据库:使用DTS(数据传输服务)等托管工具,支持并行迁移。
并行迁移(提高吞吐量)
对于超大表(如10亿条),单线程太慢,可以将ID范围拆成多个独立的区间,开多个进程/线程并行处理。
思路:
- 获取
MIN(id)和MAX(id)。 - 计算总数据量,按N个并发任务,将ID范围分成N等份(假设ID相对均匀)。
- 每个线程独立执行方案一或方案二,处理自己的ID段。
注意并发控制:
- 目标端:批量写入时,使用
INSERT ... ON DUPLICATE KEY UPDATE(防止重复)或者分批提交。 - 源端:读操作天然支持并发,但要注意数据库连接数限制。
使用ETL工具(省心但需运维)
如果不想写大量代码,可以使用成熟工具:
- DataX (阿里):支持多种数据源,自带分片算法。
- Kettle (Pentaho):图形化界面,支持分页和记录集限制。
- Apache Flink / Spark:处理实时或历史数据迁移,适合复杂转换。
关键注意事项
- 内存管理:绝对不要
fetchall()一次拉取所有数据,使用fetchmany(size)或逐行处理。 - 事务控制:
- 源端:如果迁移过程中源表还在写,需要考虑使用
SELECT ... FOR UPDATE(影响性能)或使用数据库的Binlog/Flink CDC进行实时同步。 - 目标端:建议每批提交一次(
batch_commit),而不是每行或最后一次性提交,每批提交后,如果进程崩溃,最多丢失当前这一批数据(可重跑)。
- 源端:如果迁移过程中源表还在写,需要考虑使用
- 断点续传:
- 方案一:记录
last_processed_id。 - 方案二:记录上次游标值。
- 建议将进度写入到一个持久化位置(如数据库表、Redis、文件),以便失败重启时跳过已完成的批次。
- 方案一:记录
- 性能监控:
- 开启慢查询日志。
- 监控源库的
IOPS、CPU,不要在业务高峰期运行。 - 分批写入目标时,适当调整
BATCH_SIZE(通常500-5000条/批,或每批数据量不超过1MB),避免目标库日志暴增或锁表。
- 数据类型:
- 注意大字段(
TEXT,BLOB),如果包含大字段,批次应该更小(如200条),否则可能内存溢出或网络超时。 - 注意时区、字符集转换。
- 注意大字段(
- 网络:如果源和目标在不同网络/机房,考虑使用压缩传输。
- 目标库约束:迁移前禁用目标表的外键检查和唯一性检查(如
SET FOREIGN_KEY_CHECKS=0; SET UNIQUE_CHECKS=0;),完成后再启用,可以大幅提高插入速度。
总结选择建议
- 表有自增ID,数据量大,网络好 → 方案一(范围分片) + 并行迁移。
- 表没有自增ID,但有其他有序字段 → 方案二(游标分页)。
- 纯数据库到数据库,且在同一网络 → 考虑
COPY/LOAD DATA方案。 - 需要转换数据结构(如MySQL到MongoDB) → 使用 DataX 或 Apache Spark。
- 需要不停机迁移 → 结合 Binlog (CDC) + 数据全量迁移 + 增量追同步(这是更高级的方案,但正确性最高)。
务必先在测试环境验证:用生产数据的1/100测试一下,确认数据完整性、迁移速度、内存消耗,然后再运行。