refactor: 备份/恢复改为「SQLite 归档 + 单事务整库替换」——不再换文件,也不再依赖 SQLite 运行时

SQLite 从"运行时数据库"降级为"备份交换格式":导出把当前库(MySQL 或 SQLite)
整库写成一份 SQLite 归档,恢复则是启动时在单个事务里 DELETE + INSERT。
前端三个接口的契约一行未改。

- core/system_backup.py 重写:
  * dump_to_sqlite()(旧名 snapshot_db 保留为别名)取代在线备份 API:MySQL 下把
    这条连接提到 REPEATABLE READ,保证 12 张表读的是同一时刻
  * zip 结构不变(users.db + manifest.json + apks/),既有老备份继续可导入;不依赖
    任何外部二进制(sqlite3 是标准库,python:slim 容器里现成)
  * manifest 增加 db_backend / deployment_env / source_db_id,用于识别备份来源环境
  * apply_restore:安全网从裸 pre_restore_*.db 升级为 pre_restore_*.zip(可直接再导入
    回滚);跨环境导入默认硬停(需 force_env_mismatch=true)
  * consume_pending_restore:启动时单事务整库替换,任何一步失败 rollback,当前数据
    完好(比换文件更安全);坏归档挪 restore_failed_* 且不阻塞启动
  * 归档库收尾时 checkpoint 回单文件模式并清掉 -wal/-shm(zip 只打包主文件)
- web_server.py:consume 从 init_db 之前移到之后(现在是事务替换,需要表与 app context)
- config.py:BACKUP_DIR / RESTORE_STAGING_DIR / RESTORE_PENDING_DIR 支持环境变量覆盖
  —— 自动化测试必须把恢复目录指到临时位置,否则测试造的"待生效恢复任务"会在服务
  下次重启时被当成用户的操作消费掉
- core/db_config.py:SQLite 回退模式也写库环境标签(备份要能标出来源环境)
- web/system_api.py + static/admin/system.js + templates/admin/monitor.html:
  预览显示来源/当前环境,跨环境时出现「允许跨环境导入」勾选项;文案同步

验证(隔离副本,绝不碰真实库):导出→zip 结构与 manifest 正确;预览带来源环境与
告警;应用→生成 pre_restore zip + 归档就位;跨环境 apply 被拒;模拟重启→数据回到
导出时刻、经验库/动作库完好;塞入坏归档→挪 restore_failed_*、服务照常启动、数据未变。
This commit is contained in:
2026-09-13 10:34:58 +08:00
parent b145a007eb
commit 4899e66c68
7 changed files with 321 additions and 112 deletions
+273 -90
View File
@@ -1,18 +1,25 @@
"""系统数据备份导出/导入(数据库 + APK 文件)。
导出:用 sqlite3 在线备份 API 对 data/users.db 做一致性快照 → zip
(users.db + manifest.json + 可选 apks/*.apk)。
导入:上传(zip/db) → 暂存校验预览 → 应用(先把当前库快照到
BACKUP_DIR/pre_restore_*.db,再把暂存库移到 RESTORE_PENDING_DIR)→
下次启动 consume_pending_restore() 换位生效。
**归档格式:zip(`users.db` + `manifest.json` + 可选 `apks/*.apk`)。**
SQLite 在这里的角色是**交换格式**,不是运行时数据库——这在迁到 MySQL 之后依然成立,
带来三个好处:
1. 前端三个接口的请求/响应结构一行都不用改
2. 校验逻辑(完整性 / 必需表 / schema 版本 / 未登记表反向自检)可以整体沿用
3. 2026-09-13 之前导出的老备份继续可导入;也不依赖任何外部二进制
(`sqlite3` 是 Python 标准库,python:slim 容器里现成)
为什么必须重启生效:web_server 单进程内 SQLAlchemy engine + 多个常驻后台线程
(TaskManager 内存态 / 设备发现 / Agent 等)持有 users.db(WAL);Windows 下
不能直接替换正被打开的文件,且 TaskManager 启动时才把 groups/jobs 读入内存。
故导入只落一个「待生效恢复任务」,由 web_server.py 在 init_db 之前消费。
导出:把**当前库**(MySQL 或 SQLite)整库读出来,写进一个临时 SQLite 文件 → zip。
MySQL 下把这条连接提到 REPEATABLE READ,保证所有表读的是同一时刻。
导入:上传(zip/db) → 暂存校验预览 → 应用(先自动导出一份当前库到
BACKUP_DIR/pre_restore_*.zip 作安全网,再把暂存归档移到 RESTORE_PENDING_DIR)
→ 下次启动 consume_pending_restore() 在**单个事务内**整库替换生效。
安全红线:全程只读现有库 / 新增或移动自己目录下的文件;绝不 kill-server、
绝不 disconnect(与全项目一致)。
为什么仍然要重启生效:TaskManager / device_pool 把任务与设备**缓存在内存**里,
整库替换后内存副本全部失效。(SQLite 时代还有"Windows 无法替换被持有的文件"这层,
改用 MySQL 之后只剩内存态这一层。)
安全红线:全程只读现有库 / 只写自己目录下的文件;替换数据用单事务,失败回滚,
绝不留下半新半旧的库;绝不 kill-server、绝不 disconnect(与全项目一致)。
"""
import json
import os
@@ -23,6 +30,8 @@ import uuid
import zipfile
from urllib.parse import quote
from sqlalchemy import create_engine, inspect, select, text
from config import (DATA_DIR, APK_DIR, BACKUP_DIR,
RESTORE_STAGING_DIR, RESTORE_PENDING_DIR)
from core.logger import get_logger
@@ -59,9 +68,19 @@ def _ts():
return time.strftime("%Y%m%d_%H%M%S")
def _sqlite_uri(path):
"""把路径转成 sqlite URI(兼容含空格/中文/反斜杠的 Windows 路径)。"""
return "file:{}".format(quote(os.path.abspath(path).replace("\\", "/"), safe="/:"))
def _readonly_uri(path):
"""把文件路径转成 sqlite 只读 URI(兼容含空格/中文/反斜杠的 Windows 路径)。"""
return "file:{}?mode=ro".format(quote(os.path.abspath(path).replace("\\", "/"), safe="/:"))
"""sqlite 只读 URI。"""
return _sqlite_uri(path) + "?mode=ro"
def _sqlite_engine_url(path):
"""SQLAlchemy 用的 SQLite URL(注意与 file: URI 不同,engine 不吃那个格式)。"""
return "sqlite:///" + os.path.abspath(path).replace("\\", "/")
def _connect_readonly(path):
@@ -107,34 +126,6 @@ def _table_rows(con, table):
return None
# ================== 在线快照 ==================
def snapshot_db(dest_path, src_path=None):
"""对 src(默认当前 users.db)做 sqlite 在线一致快照到 dest_path。
用标准库备份 API(不依赖 SQLAlchemy engine,WAL 模式下读一致快照安全)。
备份完成后把目标强制落回单文件(journal_mode=DELETE + checkpoint),
避免副本以 WAL 模式残留 -wal/-shm 或主文件缺刚提交帧。
"""
src_path = src_path or DB_FILE
if not os.path.exists(src_path):
raise BackupError(f"数据库不存在: {src_path}")
con = _connect_readonly(src_path)
out = sqlite3.connect(dest_path)
try:
con.backup(out)
try:
out.execute("PRAGMA journal_mode=DELETE") # checkpoint 并转回 DELETE
except sqlite3.Error:
pass
finally:
out.close()
con.close()
# 防御:确保没有任何残留附属文件
remove_quiet(dest_path + "-wal")
remove_quiet(dest_path + "-shm")
# ================== 导出 ==================
def _summary_info(con):
"""表行数摘要(只列出存在的表)。"""
rows = []
@@ -161,13 +152,123 @@ def _prune_old_exports(max_age=3600):
pass
def _prune_stale_staging():
"""清理超时未应用的暂存目录(TTL 后自动删除)。"""
now = time.time()
if not os.path.isdir(RESTORE_STAGING_DIR):
return
for name in os.listdir(RESTORE_STAGING_DIR):
p = os.path.join(RESTORE_STAGING_DIR, name)
try:
if os.path.isdir(p) and now - os.path.getmtime(p) > _STAGE_TTL:
remove_quiet(p)
except OSError:
pass
def _current_env_label():
"""当前库登记的环境标签(dev/prod)与库唯一标识(未登记则空)。"""
from core import db_config
try:
return db_config.meta_get("deployment_env") or "", db_config.meta_get("deployment_id") or ""
except Exception:
return "", ""
# ================== 导出:当前库 → SQLite 归档 ==================
def dump_to_sqlite(dest_path):
"""把当前数据库整库导出成一个 SQLite 文件(备份归档格式)。
一致性:MySQL 下把这条连接提到 REPEATABLE READ,让所有表读同一个时间点
(应用全局是 READ COMMITTED,事务提交后自增/被改的行会"漂移")。
SQLite 下不需要——单文件库天然一致。
"""
_ensure_dir(os.path.dirname(os.path.abspath(dest_path)))
remove_quiet(dest_path)
remove_quiet(dest_path + "-wal")
remove_quiet(dest_path + "-shm")
out_eng = create_engine(_sqlite_engine_url(dest_path))
is_mysql = db.engine.dialect.name == "mysql"
conn = db.engine.connect()
if is_mysql:
conn = conn.execution_options(isolation_level="REPEATABLE READ")
dumped = []
try:
# 归档库用同一套模型建表(mysql_* 表选项在 SQLite 下被忽略)
db.metadata.create_all(out_eng)
src_cols = {}
insp = inspect(conn)
for table in db.metadata.sorted_tables:
try:
src_cols[table.name] = {c["name"] for c in insp.get_columns(table.name)}
except Exception:
src_cols[table.name] = set()
with out_eng.begin() as out:
for table in db.metadata.sorted_tables:
have = src_cols.get(table.name) or set()
cols = [c for c in table.columns if c.name in have]
if not cols:
_log.warning(f"导出跳过 {table.name}:源库没有这张表")
continue
try:
rows = [dict(r._mapping)
for r in conn.execute(select(*cols))]
except Exception as e:
_log.error(f"导出读取 {table.name} 失败: {e}")
raise BackupError(f"导出失败:读取 {table.name} 出错({e})")
if rows:
out.execute(table.insert(), rows)
dumped.append((table.name, len(rows)))
finally:
conn.close()
out_eng.dispose()
# 归档库收尾:落回单文件 + 补"空值不参与唯一"的唯一索引。
# SQLAlchemy 建的 sqlite 连接会被 PRAGMA 监听器设成 WAL,留下 -wal/-shm;
# zip 里只打包主文件,所以必须先 checkpoint 回 DELETE 模式再删附属文件。
try:
con = sqlite3.connect(dest_path)
try:
con.execute("PRAGMA journal_mode=DELETE")
except sqlite3.Error as e:
_log.warning(f"归档库切回单文件模式失败: {e}")
for name, sql in (("ux_device_name",
"CREATE UNIQUE INDEX IF NOT EXISTS ux_device_name "
"ON device(name) WHERE name IS NOT NULL AND name <> ''"),
("ux_device_fingerprint",
"CREATE UNIQUE INDEX IF NOT EXISTS ux_device_fingerprint "
"ON device(fingerprint) WHERE fingerprint IS NOT NULL "
"AND fingerprint <> ''")):
try:
con.execute(sql)
except Exception as e:
_log.warning(f"归档库建索引 {name} 失败(数据可能有重复): {e}")
con.commit()
con.close()
except Exception as e:
_log.warning(f"归档库收尾失败(不影响导出): {e}")
remove_quiet(dest_path + "-wal")
remove_quiet(dest_path + "-shm")
_log.info("已导出到 SQLite 归档: %s(%d 张表)",
os.path.basename(dest_path), len(dumped))
return dumped
# 兼容旧名(历史代码/脚本里叫 snapshot_db)
snapshot_db = dump_to_sqlite
def create_export(include_apk=True):
"""生成导出 zip,返回 (zip_path, filename, manifest)。"""
_ensure_dir(BACKUP_DIR)
_prune_old_exports()
base = os.path.join(BACKUP_DIR, f"export_{_ts()}")
snap = base + ".snapshot.db"
snapshot_db(snap)
dump_to_sqlite(snap)
apk_meta = []
if include_apk and os.path.isdir(APK_DIR):
@@ -179,14 +280,18 @@ def create_export(include_apk=True):
except OSError:
pass
env_label, db_id = _current_env_label()
zip_path = base + ".zip"
con = _connect_readonly(snap)
try:
manifest = {
"format": "auto_control_backup",
"version": 1,
"version": 2,
"created_at": time.strftime("%Y-%m-%d %H:%M:%S"),
"schema_version": _read_schema_version(con),
"db_backend": db.engine.dialect.name, # mysql / sqlite
"deployment_env": env_label, # 来源环境(跨环境导入时告警)
"source_db_id": db_id, # 来源库唯一标识
"include_apk": include_apk,
"tables": _summary_info(con),
"apks": apk_meta,
@@ -237,7 +342,7 @@ def validate_backup(db_file):
if schema_version < CURRENT_SCHEMA_VERSION:
warnings.append(
f"备份 schema 较旧(v{schema_version} < 当前 v{CURRENT_SCHEMA_VERSION}),"
f"应用后首次启动会自动迁移补齐结构与新表")
f"缺失的表/列按当前结构补空,不会报错但那些数据回不来")
elif schema_version > CURRENT_SCHEMA_VERSION:
warnings.append(
f"备份 schema 较新(v{schema_version} > 当前 v{CURRENT_SCHEMA_VERSION}),"
@@ -245,14 +350,15 @@ def validate_backup(db_file):
if missing_optional:
warnings.append("备份缺少部分可选表(" + "、".join(
TABLE_LABELS.get(t, t) for t in missing_optional)
+ "),应用后启动会自动补建空表")
+ "),导入后这些表为空")
# 反向自检:备份里有"未登记"的表 → 提醒把它纳入覆盖清单(红线)
extra_tables = sorted(t for t in tables
if t not in SUMMARY_TABLES and not t.startswith("sqlite_"))
if extra_tables:
warnings.append("备份含未登记的其它表(" + "、".join(extra_tables)
+ ")——如属业务数据,请登记进 core/system_backup.py 的覆盖清单")
warnings.append("备份为全量快照:含用户口令哈希、AI 配置里的 API Key 等敏感信息,请妥善保管")
+ ")——如属业务数据,请登记进 core/models.py 的模型与 "
"core/system_backup.py 的 TABLE_LABELS")
warnings.append("备份为全量数据:含用户口令哈希、AI 配置里的 API Key 等敏感信息,请妥善保管")
return True, {
"integrity": integrity,
"schema_version": schema_version,
@@ -270,7 +376,7 @@ def validate_backup(db_file):
# ================== 导入:暂存 + 预览 ==================
def _extract_zip(zip_path, staging_dir):
"""解压备份 zip:取顶层 users.db 与可选 apks/*.apk(防路径穿越,只用 basename)。"""
"""解压备份 zip:取顶层 users.db、可选 apks/*.apk 与 manifest.json(防路径穿越)。"""
db_dest = os.path.join(staging_dir, "users.db")
apk_dest = os.path.join(staging_dir, "apks")
got_db = False
@@ -283,6 +389,10 @@ def _extract_zip(zip_path, staging_dir):
if name == "users.db" and base == "users.db":
z.extract(info, staging_dir) # 已确认顶层名字,无穿越
got_db = True
elif name == "manifest.json":
with z.open(info) as src, open(
os.path.join(staging_dir, "manifest.json"), "wb") as dst:
shutil.copyfileobj(src, dst)
elif name.startswith("apks/") and base.lower().endswith(".apk"):
_ensure_dir(apk_dest)
with z.open(info) as src, open(
@@ -293,18 +403,15 @@ def _extract_zip(zip_path, staging_dir):
return db_dest
def _prune_stale_staging():
"""清理超时未应用的暂存目录(TTL 后自动删除)。"""
now = time.time()
if not os.path.isdir(RESTORE_STAGING_DIR):
return
for name in os.listdir(RESTORE_STAGING_DIR):
p = os.path.join(RESTORE_STAGING_DIR, name)
try:
if os.path.isdir(p) and now - os.path.getmtime(p) > _STAGE_TTL:
remove_quiet(p)
except OSError:
pass
def _read_staged_manifest(staging_dir):
p = os.path.join(staging_dir, "manifest.json")
if not os.path.exists(p):
return {}
try:
with open(p, encoding="utf-8") as f:
return json.load(f) or {}
except Exception:
return {}
def stage_upload(file_storage):
@@ -326,12 +433,25 @@ def stage_upload(file_storage):
db_file = os.path.join(staging_dir, "users.db")
os.replace(raw_path, db_file)
else:
raise BackupError("仅支持 .zip(平台导出)或 .db(sqlite 库)文件")
raise BackupError("仅支持 .zip(平台导出)或 .db(备份库文件)")
ok, info = validate_backup(db_file)
if not ok:
raise BackupError(info)
info["file_name"] = orig_name
info["file_size"] = os.path.getsize(db_file)
# 跨环境提示(防混库第 5 层):备份来自哪个环境、当前连的是哪个环境
man = _read_staged_manifest(staging_dir)
info["source_env"] = man.get("deployment_env") or ""
info["source_db_id"] = man.get("source_db_id") or ""
info["source_backend"] = man.get("db_backend") or "sqlite(旧版备份)"
cur_env, _ = _current_env_label()
info["current_env"] = cur_env
if info["source_env"] and cur_env and info["source_env"] != cur_env:
info["warnings"].append(
f"该备份来自 {info['source_env']} 环境,而当前连接的是 {cur_env} 库 —— "
f"导入会覆盖当前全部数据。确认这是你要的操作再继续(跨环境导入默认被拒绝)。")
# 暂存目录里留住 manifest,apply 时还要用它做环境校验
return token, info
except BackupError:
remove_quiet(staging_dir)
@@ -342,8 +462,8 @@ def stage_upload(file_storage):
# ================== 导入:应用(落待生效任务) ==================
def apply_restore(token):
"""校验 token → 快照当前库 → 把暂存库移到 restore_pending。
def apply_restore(token, force_env_mismatch=False):
"""校验 token → 自动备份当前库(zip 安全网)→ 把暂存归档移到 restore_pending。
返回 dict {ok, backup_name, message};失败抛 BackupError。token 形如 12 位 hex。
"""
@@ -357,25 +477,41 @@ def apply_restore(token):
if not ok:
raise BackupError(info)
# 1) 自动备份当前库(安全网,可回滚)
# 跨环境硬停:默认不许把 prod 的备份灌进 dev 库(反之亦然)
man = _read_staged_manifest(staging_dir)
src_env = man.get("deployment_env") or ""
cur_env, _ = _current_env_label()
if src_env and cur_env and src_env != cur_env and not force_env_mismatch:
raise BackupError(
f"备份来自 {src_env} 环境,当前库是 {cur_env} 环境 —— 已拒绝导入。\n"
f"如确需跨环境恢复,请在确认页勾选「允许跨环境导入」后重试。")
# 1) 自动备份当前库(安全网,可回滚;导出成 zip,能直接再导入回来)
_ensure_dir(BACKUP_DIR)
backup_name = f"pre_restore_{_ts()}.db"
snapshot_db(os.path.join(BACKUP_DIR, backup_name))
try:
zip_path, _, _ = create_export(include_apk=True)
backup_name = f"pre_restore_{_ts()}.zip"
shutil.move(zip_path, os.path.join(BACKUP_DIR, backup_name))
except Exception as e:
raise BackupError(f"导入前自动备份当前库失败,已中止导入: {e}")
# 2) 落待生效恢复任务(旧未消费任务被覆盖——上一份已无意义)
pending_dir = RESTORE_PENDING_DIR
if os.path.isdir(pending_dir):
remove_quiet(pending_dir)
_ensure_dir(pending_dir)
os.replace(db_file, os.path.join(pending_dir, "users.db"))
if os.path.isdir(RESTORE_PENDING_DIR):
remove_quiet(RESTORE_PENDING_DIR)
_ensure_dir(RESTORE_PENDING_DIR)
os.replace(db_file, os.path.join(RESTORE_PENDING_DIR, "users.db"))
for extra in ("manifest.json",):
src = os.path.join(staging_dir, extra)
if os.path.exists(src):
shutil.move(src, os.path.join(RESTORE_PENDING_DIR, extra))
apk_src = os.path.join(staging_dir, "apks")
if os.path.isdir(apk_src) and os.listdir(apk_src):
shutil.move(apk_src, os.path.join(pending_dir, "apks"))
shutil.move(apk_src, os.path.join(RESTORE_PENDING_DIR, "apks"))
remove_quiet(staging_dir)
schema_version = info.get("schema_version", 0)
_log.info(f"导入已落恢复任务: schema v{schema_version},"
f"重启 web_server 后生效;当前库已快照 {backup_name}")
f"重启 web_server 后生效;当前库已自动备份 {backup_name}")
return {
"ok": True,
"backup_name": backup_name,
@@ -383,32 +519,43 @@ def apply_restore(token):
}
# ================== 启动消费(web_server.py init_db 前调用) ==================
def consume_pending_restore():
"""把 restore_pending/users.db 换位为当前 users.db(删除旧 wal/shm,合并 apks)。
# ================== 启动消费(init_db 之后、TaskManager 之前调用)==================
def _pending_db_path():
"""待生效恢复的归档路径(兼容旧版扁平布局)。"""
p = os.path.join(RESTORE_PENDING_DIR, "users.db")
return p if os.path.exists(p) else None
必须在 SQLAlchemy engine 首次打开 users.db 之前执行(web_server 装配时调用)。
失败不阻塞启动:坏恢复文件会挪到 BACKUP_DIR/restore_failed_*/ 并继续用旧库。
def consume_pending_restore():
"""把 restore_pending/users.db 的内容在**单个事务内**整库替换进来。
与 SQLite 时代的"换文件"不同:现在是 DELETE 全表 + 分块 INSERT,全程一个事务——
任何一步失败就 rollback,当前数据保持原样(比换文件更安全)。
必须在 init_db() 之后调用(表已建好、app context 可用);失败不阻塞启动:
坏归档挪到 BACKUP_DIR/restore_failed_*/ 并继续用当前数据。
返回是否真的应用了恢复。
"""
pending_db = os.path.join(RESTORE_PENDING_DIR, "users.db")
if not os.path.exists(pending_db):
pending_db = _pending_db_path()
if not pending_db:
return False
# 消费前轻量复验(防止意外损坏文件把主库换掉)
# 消费前轻量复验(防止意外损坏的文件把主库换掉)
ok, msg = validate_backup(pending_db)
if not ok:
fail_dir = os.path.join(BACKUP_DIR, f"restore_failed_{_ts()}")
_ensure_dir(fail_dir)
os.replace(pending_db, os.path.join(fail_dir, "users.db"))
shutil.move(pending_db, os.path.join(fail_dir, "users.db"))
remove_quiet(RESTORE_PENDING_DIR)
_log.error(f"恢复任务校验失败,已搁置到 {fail_dir},继续使用当前库: {msg}")
return False
_ensure_dir(DATA_DIR)
os.replace(pending_db, DB_FILE)
for ext in ("-wal", "-shm"):
p = DB_FILE + ext
if os.path.exists(p):
remove_quiet(p)
try:
applied = _apply_archive_to_db(pending_db)
except Exception as e:
_log.error(f"应用备份恢复失败,已回滚(当前数据未改动): {e},"
f"归档保留在 {RESTORE_PENDING_DIR} 供排查")
return False
# 合并 apks(覆盖同名)
apk_src = os.path.join(RESTORE_PENDING_DIR, "apks")
if os.path.isdir(apk_src):
@@ -420,6 +567,42 @@ def consume_pending_restore():
os.remove(dst)
shutil.move(src, dst)
remove_quiet(RESTORE_PENDING_DIR)
schema_version = msg.get("schema_version", 0) if isinstance(msg, dict) else "?"
_log.info(f"已应用备份恢复: users.db 替换完成(schema v{schema_version}),apks 已合并")
schema_version = msg.get("schema_version", "?") if isinstance(msg, dict) else "?"
_log.info(f"已应用备份恢复(单事务替换): {applied},schema v{schema_version},apks 已合并")
return True
def _apply_archive_to_db(db_file, chunk=500):
"""把 SQLite 归档里的数据整库写进当前数据库,单事务、失败回滚。
用 DELETE 而不是 TRUNCATE:TRUNCATE 是 DDL,会隐式提交,破坏原子性。
"""
con = _connect_readonly(db_file)
arch_tables = _list_tables(con)
stats = []
try:
# 单事务:MySQL/SQLite 都支持事务内 DELETE + INSERT
with db.engine.begin() as dst:
for table in reversed(db.metadata.sorted_tables):
dst.execute(table.delete())
for table in db.metadata.sorted_tables:
if table.name not in arch_tables:
_log.warning(f"归档里没有 {table.name},该表导入后为空")
continue
have = {r[1] for r in con.execute(
'PRAGMA table_info("%s")' % table.name)}
cols = [c for c in table.columns if c.name in have]
if not cols:
_log.warning(f"归档表 {table.name} 没有可用列,跳过")
continue
sel = "SELECT {} FROM \"{}\"".format(
", ".join('"%s"' % c.name for c in cols), table.name)
rows = con.execute(sel).fetchall()
names = [c.name for c in cols]
for i in range(0, len(rows), chunk):
batch = [dict(zip(names, r)) for r in rows[i:i + chunk]]
dst.execute(table.insert(), batch)
stats.append(f"{table.name}={len(rows)}")
finally:
con.close()
return " ".join(stats)