怎样用脚本批量同步数据库数据?

wen 实用脚本 1

本文目录导读:

怎样用脚本批量同步数据库数据?

  1. 核心思路
  2. 同构数据库(MySQL → MySQL)
  3. 异构数据库(MySQL → PostgreSQL / SQL Server)
  4. ETL 工具 vs 脚本
  5. 注意事项

批量同步数据库数据通常需要根据数据量、源库与目标库的形态(同构/异构)、同步频率等选择合适的脚本方案,下面给出几种常见场景的脚本实现思路与示例。


核心思路

  1. 全量同步:清空目标表,重新插入源库所有数据。
  2. 增量同步:基于时间戳、自增ID、变更日志(CDC)或触发器,仅同步变化的数据。
  3. 变更捕获方式
    • 时间戳字段(updated_at)
    • 自增ID 断点(记录上次同步的最大ID)
    • 对比MD5(适用于少量字段、低频场景)
  4. 执行机制: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、断网恢复

注意事项

  1. 事务与一致性:大表同步建议分页 + 快照读(SET TRANSACTION ISOLATION LEVEL REPEATABLE READ)。
  2. 主键冲突:使用 INSERT ... ON DUPLICATE KEY UPDATE(MySQL)或 ON CONFLICT(PG)。
  3. 时区:确保源库与目标库 time_zone 一致,或显式转换。
  4. 监控:记录每次同步的 start_time / end_time / row_count 到监控表或日志。
  5. 错误重试:连接丢失、死锁等异常应重试3次,间隔指数退避。

如果需要针对 具体数据库类型(Oracle → MySQL、MongoDB → PostgreSQL)或 特定同步频率(实时、每小时)的脚本,可以进一步说明,我可以给出更详细的示例。

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