Files
BLD_sync/run_incremental_sync.py
Misaka_Company 3791966446 Improve transaction management and connection cleanup in incremental sync
- Add explicit BEGIN TRANSACTION for clear transaction boundaries
- Ensure IDENTITY_INSERT state is properly tracked and cleaned up
- Add finally block to guarantee Access connection cleanup
- Move table existence check inside sync loop for better scope

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-06-11 15:33:57 +08:00

286 lines
12 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.
import time
import os
import sys
import pyodbc
from config import SQL_SERVER_CONN, ACCESS_DRIVER, SYNC_MAPPING, POLL_INTERVAL, LOG_TABLE_CONFIG, UPTIME_KUMA_CONFIG
import db_utils
from log_utils import (LoggerManager, log_success, log_error, log_warning, log_info, log_processing,
log_skip, log_critical, log_start, log_stop, log_file, log_database, log_sync)
from uptime_kuma_utils import UptimeKumaMonitor
# 初始化日志管理器
LoggerManager("run_incremental_sync", log_prefix="incremental")
# 初始化 Uptime Kuma 监控器
uptime_monitor = UptimeKumaMonitor(UPTIME_KUMA_CONFIG)
uptime_monitor.set_logger(log_warning)
# ================= 辅助函数 =================
def has_identity_column(sql_cursor, schema, table):
"""检查表是否包含标识列"""
query = """
SELECT COUNT(*)
FROM sys.columns c
JOIN sys.tables t ON c.object_id = t.object_id
JOIN sys.schemas s ON t.schema_id = s.schema_id
WHERE s.name = ? AND t.name = ? AND c.is_identity = 1
"""
sql_cursor.execute(query, (schema, table))
return sql_cursor.fetchone()[0] > 0
# ================= 主逻辑 =================
def process_sync_task():
"""
逻辑重构:
遍历 SYNC_MAPPING 中的每一个文件 -> 去日志表查询该文件下特定表的未同步记录。
"""
sql_conn = db_utils.get_sql_conn()
sql_cursor = sql_conn.cursor()
sql_cursor.fast_executemany = True # 高性能开关
# 标记是否有工作被处理(用于控制轮询休眠时间)
work_done = False
cols = LOG_TABLE_CONFIG
log_full_name = db_utils.fmt_table(cols['schema'], cols['table_name'])
try:
# === 核心循环:以配置文件为驱动 ===
for clean_path, tables_map in SYNC_MAPPING.items():
# 1. 准备查询条件
# 获取该文件下所有需要同步的表名列表
target_tables = list(tables_map.keys())
if not target_tables:
continue
# 构造 TableAddress 的精确匹配条件
# VBA 逻辑:网络路径带 ";DATABASE=", 本地路径带 "LOCAL:"
# 我们直接构造这两个字符串,让 SQL Server 做精确匹配,效率极高
# 注意:将路径转为 Windows 标准反斜杠
win_path = os.path.normpath(clean_path)
addr_candidates = [
f";DATABASE={win_path}", # 情况1
f"LOCAL={win_path}", # 情况2 (注意 VBA 代码里可能是 LOCAL: 或 LOCAL=,请核对)
f"LOCAL:{win_path}", # 情况3
win_path # 情况4 (兼容没有前缀的情况)
]
# 2. 构造动态 SQL 查询
# WHERE Synced=0 AND Address IN (...) AND TableName IN (...)
placeholders_addr = ','.join(['?'] * len(addr_candidates))
placeholders_tbl = ','.join(['?'] * len(target_tables))
query_log = f"""
SELECT TOP 1000
{cols['col_log_id']},
{cols['col_table_name']},
{cols['col_record_id']}
FROM {log_full_name}
WHERE {cols['col_synced']} = 0
AND TableType = 'LINKED_ACCESS'
AND {cols['col_address']} IN ({placeholders_addr})
AND {cols['col_table_name']} IN ({placeholders_tbl})
ORDER BY {cols['col_log_id']} ASC
"""
# 参数列表:先放地址,再放表名
params = addr_candidates + target_tables
sql_cursor.execute(query_log, params)
logs = sql_cursor.fetchall()
if not logs:
continue # 这个文件没有需要同步的记录,检查下一个文件
work_done = True # 标记有工作
log_file(f"{os.path.basename(clean_path)} 发现 {len(logs)} 条待同步变更")
# 3. 本地分组 (按表名)
# 结构: table_tasks[TableName] = { ids: {}, log_ids: [] }
table_tasks = {}
for row in logs:
log_id, acc_table, record_id = row
if acc_table not in table_tasks:
table_tasks[acc_table] = {'record_ids': set(), 'log_ids': []}
table_tasks[acc_table]['record_ids'].add(record_id)
table_tasks[acc_table]['log_ids'].append(log_id)
# 4. 执行同步 (连接一次 Access处理多张表)
if not os.path.exists(clean_path):
log_error(f"无法访问文件: {clean_path}")
continue
acc_conn = None
try:
acc_conn = db_utils.get_access_conn(clean_path)
acc_cursor = acc_conn.cursor()
log_database(f"已连接 Access 文件: {os.path.basename(clean_path)}")
except Exception as conn_err:
log_error(f"连接 Access 失败 [{os.path.basename(clean_path)}]: {conn_err}")
continue
# 统计每个文件的同步情况
file_success_count = 0
file_error_count = 0
file_total_records = 0
for acc_table, data in table_tasks.items():
record_ids = list(data['record_ids'])
log_ids = data['log_ids']
# 读取目标配置
target_conf = tables_map[acc_table]
target_schema = target_conf['target_schema']
target_table = target_conf['target_table']
pk_col = target_conf['pk_col']
target_full_name = db_utils.fmt_table(target_schema, target_table)
log_processing(f"正在同步表 [{acc_table}] → [{target_table}] ({len(record_ids)} 条记录)")
# 检查目标表是否存在,不存在则自动创建
if not db_utils.table_exists(sql_cursor, target_schema, target_table):
db_utils.ensure_schema(sql_cursor, target_schema)
acc_cursor.execute(f"SELECT TOP 1 * FROM [{acc_table}]")
db_utils.create_table_from_access(
sql_cursor, target_schema, target_table,
acc_cursor.description, pk_col
)
log_info(f"目标表 [{target_table}] 不存在,已自动创建")
# 事务状态追踪
identity_enabled = False
transaction_begun = False
try:
# --- 开始事务 ---
sql_cursor.execute("BEGIN TRANSACTION")
transaction_begun = True
# --- A. Access 查新数据 ---
ids_placeholders = ','.join(['?'] * len(record_ids))
acc_sql = f"SELECT * FROM [{acc_table}] WHERE [{pk_col}] IN ({ids_placeholders})"
acc_cursor.execute(acc_sql, record_ids)
new_rows = acc_cursor.fetchall()
acc_cols = [col[0] for col in acc_cursor.description]
# --- B. SQL Server 删旧插新 ---
# 1. 删除旧数据
del_sql = f"DELETE FROM {target_full_name} WHERE [{pk_col}] IN ({ids_placeholders})"
sql_cursor.execute(del_sql, record_ids)
# 2. 插入新数据
if new_rows:
insert_sql = db_utils.generate_insert_sql(target_schema, target_table, acc_cols)
# 检查并启用 IDENTITY_INSERT
has_identity = has_identity_column(sql_cursor, target_schema, target_table)
if has_identity:
sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} ON")
identity_enabled = True
sql_cursor.executemany(insert_sql, new_rows)
# 关闭 IDENTITY_INSERT
if identity_enabled:
sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} OFF")
identity_enabled = False
# 3. 标记日志 Synced = 1
log_placeholders = ','.join(['?'] * len(log_ids))
update_log_sql = f"""
UPDATE {log_full_name}
SET {cols['col_synced']} = 1
WHERE {cols['col_log_id']} IN ({log_placeholders})
"""
sql_cursor.execute(update_log_sql, log_ids)
# 提交事务
sql_conn.commit()
transaction_begun = False
log_info(f"表 [{target_table}] 同步完成: {len(record_ids)} 条记录")
file_success_count += 1
file_total_records += len(record_ids)
except Exception as tbl_err:
# 回滚事务
if transaction_begun:
sql_conn.rollback()
transaction_begun = False
# 清理 IDENTITY_INSERT 状态
if identity_enabled:
try:
sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} OFF")
except:
pass
log_error(f"表 [{acc_table}] 同步失败: {tbl_err}")
file_error_count += 1
finally:
# 确保 Access 连接关闭
try:
if acc_conn:
acc_conn.close()
log_info(f"已关闭 Access 连接: {os.path.basename(clean_path)}")
except Exception as close_err:
log_warning(f"关闭 Access 连接时出错: {close_err}")
# 输出文件级别的汇总
if file_success_count > 0 or file_error_count > 0:
summary = f"文件 [{os.path.basename(clean_path)}] 同步汇总: "
summary += f"成功 {file_success_count} 张表 ({file_total_records} 条记录)"
if file_error_count > 0:
summary += f" | 失败 {file_error_count} 张表"
log_sync(summary)
return work_done
except Exception as e:
log_critical(f"全局异常: {e}")
return False
finally:
try: sql_conn.close()
except: pass
# ================= Uptime Kuma 心跳 =================
# 使用 uptime_kuma_utils.UptimeKumaMonitor 替代原有实现
if __name__ == "__main__":
log_start("增量同步服务已启动 (配置驱动模式)")
log_info(f"轮询间隔: {POLL_INTERVAL}")
log_info(f"监控配置: {len(SYNC_MAPPING)} 个文件")
if UPTIME_KUMA_CONFIG.get('enabled', False):
log_info(f"心跳间隔: {UPTIME_KUMA_CONFIG['heartbeat_interval']}")
log_info("=" * 70)
# 启动时发送第一次心跳
uptime_monitor.send_heartbeat()
try:
while True:
try:
has_work = process_sync_task()
# 检查是否需要发送心跳
uptime_monitor.check_and_send_heartbeat()
# 如果有工作,说明可能还有积压,休息短一点(0.1s)
# 如果没工作,休息标准间隔(5s)
time.sleep(0.1 if has_work else POLL_INTERVAL)
except KeyboardInterrupt:
log_info("=" * 70)
log_stop("收到停止信号,服务正在关闭...")
break
except Exception as e:
log_critical(f"主循环崩溃: {e}")
time.sleep(5)
finally:
# 停止时发送心跳停止信号(可选)
uptime_monitor.send_stop_signal()