Files
zentao-flow/tmp/db_copy_200_to_161.py
T

176 lines
5.9 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: utf-8 -*-
"""200 zentao_dev 只读导出 → 161 zentao_dev_2026 导入
源库只执行 SELECT/SHOW,绝不写入;流式游标 + 分批插入。
"""
import io, sys, time, re, pymysql
from pymysql.cursors import SSCursor
sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding='utf-8')
SRC = dict(host='192.168.3.200', user='root', password='PX4fTAAsJ#T!1',
database='zentao_dev', charset='utf8mb4', connect_timeout=15,
read_timeout=900, write_timeout=900)
DST = dict(host='192.168.1.161', user='devgps', password='dev@2021GPS',
charset='utf8mb4', autocommit=False, connect_timeout=15,
read_timeout=900, write_timeout=900)
TARGET_DB = 'zentao_dev_2026'
BATCH = 500
# 日志大表只建结构不拷数据(用户拍板 2026-07-29)
SKIP_DATA = {'zt_action', 'zt_actionrecent'}
def log(msg):
print(f'[{time.strftime("%H:%M:%S")}] {msg}', flush=True)
src = pymysql.connect(**SRC)
dst = pymysql.connect(**DST)
scur = src.cursor()
dcur = dst.cursor()
# 目标库会话加速+防约束干扰+放大单包上限
dcur.execute('SET SESSION foreign_key_checks=0')
dcur.execute('SET SESSION unique_checks=0')
try:
dcur.execute('SET SESSION max_allowed_packet=1073741824')
except Exception:
pass
dcur.execute(f'CREATE DATABASE IF NOT EXISTS `{TARGET_DB}` DEFAULT CHARACTER SET utf8mb4')
dcur.execute(f'USE `{TARGET_DB}`')
dst.commit()
log(f'目标库 {TARGET_DB} 就绪')
# 源库表清单(仅 BASE TABLE)
scur.execute("SELECT table_name FROM information_schema.tables WHERE table_schema='zentao_dev' AND table_type='BASE TABLE' ORDER BY table_name")
tables = [r[0] for r in scur.fetchall()]
log(f'共 {len(tables)} 张表')
# 1) 建表结构(只建目标库缺失的表;已存在的不动——配合断点续跑,已拷数据不丢)
dcur.execute("SELECT table_name FROM information_schema.tables WHERE table_schema=%s", (TARGET_DB,))
existing = {r[0] for r in dcur.fetchall()}
created = 0
for t in tables:
if t in existing:
continue
scur.execute(f'SHOW CREATE TABLE `{t}`')
create_sql = scur.fetchone()[1]
create_sql = re.sub(r' DEFINER=`[^`]+`@`[^`]+`', '', create_sql)
dcur.execute(create_sql)
created += 1
dst.commit()
log(f'表结构就绪(新建 {created} 张,复用 {len(tables)-created} 张)')
def row_size(row):
s = 0
for v in row:
if v is None:
continue
if isinstance(v, (bytes, bytearray)):
s += len(v)
elif isinstance(v, str):
s += len(v.encode('utf-8', 'ignore'))
else:
s += 16
return s
MAX_PACKET = 12 * 1024 * 1024 # 单包限 12MB(161 max_allowed_packet=64MB,留足余量)
def copy_table(t, attempt):
"""单表复制:断线重连 + 按字节分包;重试前清空目标表保证幂等"""
global src, dst, scur, dcur
if attempt > 0:
time.sleep(3)
for conn in (src, dst):
try:
conn.close()
except Exception:
pass
src = pymysql.connect(**SRC)
dst = pymysql.connect(**DST)
scur = src.cursor()
dcur = dst.cursor()
dcur.execute('SET SESSION foreign_key_checks=0')
dcur.execute('SET SESSION unique_checks=0')
try:
dcur.execute('SET SESSION max_allowed_packet=1073741824')
except Exception:
pass
dcur.execute(f'USE `{TARGET_DB}`')
dst.commit()
dcur.execute(f'DELETE FROM `{t}`')
dst.commit()
read_cur = src.cursor(SSCursor)
read_cur.execute(f'SELECT * FROM `{t}`')
cols = [d[0] for d in read_cur.description]
col_sql = ','.join(f'`{c}`' for c in cols)
ph = ','.join(['%s'] * len(cols))
insert_sql = f'INSERT INTO `{t}` ({col_sql}) VALUES ({ph})'
n = 0
batch = []
batch_bytes = 0
for row in read_cur:
batch.append(row)
batch_bytes += row_size(row)
if len(batch) >= BATCH or batch_bytes >= MAX_PACKET:
dcur.executemany(insert_sql, batch)
dst.commit()
n += len(batch)
batch = []
batch_bytes = 0
if batch:
dcur.executemany(insert_sql, batch)
dst.commit()
n += len(batch)
read_cur.close()
return n
# 2) 逐表复制数据(源只读流式;SKIP_DATA 只留空表结构;已拷完的表自动跳过=断点续跑)
total_rows = 0
for t in tables:
if t in SKIP_DATA:
log(f'{t}: 跳过数据(仅结构,0 行)')
continue
scur.execute(f'SELECT COUNT(*) FROM `{t}`')
src_cnt = scur.fetchone()[0]
dcur.execute(f'SELECT COUNT(*) FROM `{t}`')
if dcur.fetchone()[0] == src_cnt:
log(f'{t}: 已拷过({src_cnt} 行),续跑跳过')
total_rows += src_cnt
continue
t0 = time.time()
last_err = None
for attempt in range(3):
try:
n = copy_table(t, attempt)
break
except pymysql.err.OperationalError as e:
last_err = e
log(f'{t}: 第 {attempt+1} 次失败 {e},重连重试')
else:
raise last_err
total_rows += n
log(f'{t}: {n} 行 ({time.time()-t0:.0f}s)')
log(f'数据复制完成,总 {total_rows} 行')
# 3) 逐表行数对拍(SKIP_DATA 预期 0 行,单独标注)
mismatch = []
for t in tables:
scur.execute(f'SELECT COUNT(*) FROM `{t}`')
s = scur.fetchone()[0]
dcur.execute(f'SELECT COUNT(*) FROM `{t}`')
d = dcur.fetchone()[0]
if t in SKIP_DATA:
log(f'{t}: 源 {s} 行 → 目标 {d} 行(按计划跳过)')
elif s != d:
mismatch.append((t, s, d))
if mismatch:
log(f'!!! 行数不一致 {len(mismatch)} 张: {mismatch}')
else:
log(f'对拍通过:{len(tables)} 张表行数全部一致')
# 4) 抽查中文往返
dcur.execute('SELECT id, title FROM zt_story WHERE id=6566')
log(f'抽查 zt_story 6566: {dcur.fetchone()}')
dcur.execute('SELECT COUNT(*) FROM zt_perf_config')
log(f'抽查 zt_perf_config 行数: {dcur.fetchone()[0]}')
src.close(); dst.close()
log('全部完成')