本文目录导读:

批量同步数据库数据通常需要根据数据量、源库与目标库的形态(同构/异构)、同步频率等选择合适的脚本方案,下面给出几种常见场景的脚本实现思路与示例。
核心思路
- 全量同步:清空目标表,重新插入源库所有数据。
- 增量同步:基于时间戳、自增ID、变更日志(CDC)或触发器,仅同步变化的数据。
- 变更捕获方式:
- 时间戳字段(updated_at)
- 自增ID 断点(记录上次同步的最大ID)
- 对比MD5(适用于少量字段、低频场景)
- 执行机制:Shell/Python脚本驱动,支持重试、日志、异常告警。
同构数据库(MySQL → MySQL)
方案1:Shell + mysqldump(全量同步)
适合数据量小、定期全量覆盖的场景。
#!/bin/bash
# 配置
SRC_HOST="192.168.1.100"
SRC_USER="source_user"
SRC_PASS="source_pass"
SRC_DB="mydb"
DST_HOST="192.168.1.200"
DST_USER="target_user"
DST_PASS="target_pass"
DST_DB="mydb"
TABLE_LIST="users orders products"
for table in $TABLE_LIST; do
echo "同步表: $table"
mysqldump -h $SRC_HOST -u $SRC_USER -p$SRC_PASS \
--no-create-info --skip-triggers --compact $SRC_DB $table |
mysql -h $DST_HOST -u $DST_USER -p$DST_PASS $DST_DB
done
缺点:锁表风险、不适用于大表(百万级以上)。
方案2:Python + PyMySQL(增量同步 + 分页)
适合大表,支持断点续传、日志、错误处理。
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
import pymysql
import time
import logging
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
class DbSyncer:
def __init__(self):
self.src_conn = pymysql.connect(
host='192.168.1.100', user='source', password='pass', database='mydb', charset='utf8mb4'
)
self.dst_conn = pymysql.connect(
host='192.168.1.200', user='target', password='pass', database='mydb', charset='utf8mb4'
)
self.batch_size = 5000
# 断点文件,记录上次同步的最大ID或时间戳
self.checkpoint_file = '/tmp/sync_checkpoint.txt'
def read_checkpoint(self):
try:
with open(self.checkpoint_file, 'r') as f:
return int(f.read().strip())
except:
return 0
def write_checkpoint(self, value):
with open(self.checkpoint_file, 'w') as f:
f.write(str(value))
def sync_incremental(self, table, pk_column='id', time_column='updated_at'):
last_id = self.read_checkpoint()
logging.info(f"开始同步表 {table},起始ID: {last_id}")
src_cursor = self.src_conn.cursor(pymysql.cursors.DictCursor)
dst_cursor = self.dst_conn.cursor()
while True:
# 1. 从源库分页读取新增数据
sql = f"""
SELECT * FROM {table}
WHERE {pk_column} > %s
ORDER BY {pk_column} ASC
LIMIT {self.batch_size}
"""
src_cursor.execute(sql, (last_id,))
rows = src_cursor.fetchall()
if not rows:
break
# 2. 逐条写入目标库(或使用批量insert + on duplicate key update)
insert_sql = self._build_upsert_sql(table, rows[0].keys())
for row in rows:
values = [row[col] for col in row]
dst_cursor.execute(insert_sql, values)
self.dst_conn.commit()
last_id = rows[-1][pk_column]
self.write_checkpoint(last_id)
logging.info(f"表 {table} 已同步到 ID: {last_id},本次条数: {len(rows)}")
src_cursor.close()
dst_cursor.close()
logging.info(f"表 {table} 同步完成")
def _build_upsert_sql(self, table, columns):
cols = ', '.join(columns)
placeholders = ', '.join(['%s'] * len(columns))
update_part = ', '.join([f"{col}=VALUES({col})" for col in columns])
return f"""
INSERT INTO {table} ({cols}) VALUES ({placeholders})
ON DUPLICATE KEY UPDATE {update_part}
"""
def close(self):
self.src_conn.close()
self.dst_conn.close()
if __name__ == '__main__':
syncer = DbSyncer()
try:
# 支持多表顺序同步
for table in ['users', 'orders']:
syncer.sync_incremental(table, pk_column='id', time_column='updated_at')
finally:
syncer.close()
优点:
- 基于ID断点,支持增量
- 分页拉取,不撑爆内存
- UPSERT语句避免重复报错
异构数据库(MySQL → PostgreSQL / SQL Server)
推荐使用 Python + SQLAlchemy 统一接口,适配不同方言。
from sqlalchemy import create_engine, MetaData, Table, text
from sqlalchemy.dialects.postgresql import insert as pg_insert
import time
SRC_DSN = 'mysql+pymysql://user:pass@host/mydb?charset=utf8mb4'
DST_DSN = 'postgresql+psycopg2://user:pass@host/mydb'
src_engine = create_engine(SRC_DSN, pool_size=5)
dst_engine = create_engine(DST_DSN, pool_size=5)
def sync_table(table_name, chunk_size=5000):
src_conn = src_engine.connect()
dst_conn = dst_engine.connect()
# 反射表结构
metadata = MetaData()
table = Table(table_name, metadata, autoload_with=src_engine)
# 查询所有列(简单全量,可改为带断点)
query = table.select()
result = src_conn.execute(query)
while True:
rows = result.fetchmany(chunk_size)
if not rows:
break
dict_rows = [dict(row._mapping) for row in rows]
# 使用PostgreSQL的ON CONFLICT语法
stmt = pg_insert(table).values(dict_rows)
stmt = stmt.on_conflict_do_nothing()
dst_conn.execute(stmt)
src_conn.close()
dst_conn.close()
print(f"表 {table_name} 同步完成")
if __name__ == '__main__':
sync_table('users')
sync_table('orders')
ETL 工具 vs 脚本
| 场景 | 推荐方案 | 原因 |
|---|---|---|
| 简单定时全量(< 10万行) | Shell + mysqldump / pg_dump | 实现简单 |
| 增量同步、大表(百万级) | Python 脚本 + 断点 | 灵活可控 |
| 异构、复杂转换 | Python + Pandas / SQLAlchemy | 数据类型自动映射 |
| 生产级高可用 | DataX / Canal / Debezium + Kafka | 支持CDC、断网恢复 |
注意事项
- 事务与一致性:大表同步建议分页 + 快照读(SET TRANSACTION ISOLATION LEVEL REPEATABLE READ)。
- 主键冲突:使用
INSERT ... ON DUPLICATE KEY UPDATE(MySQL)或ON CONFLICT(PG)。 - 时区:确保源库与目标库 time_zone 一致,或显式转换。
- 监控:记录每次同步的 start_time / end_time / row_count 到监控表或日志。
- 错误重试:连接丢失、死锁等异常应重试3次,间隔指数退避。
如果需要针对 具体数据库类型(Oracle → MySQL、MongoDB → PostgreSQL)或 特定同步频率(实时、每小时)的脚本,可以进一步说明,我可以给出更详细的示例。