Enhance logging and synchronization features across multiple scripts
- Added detailed logging functionality to init_full_sync.py and run_incremental_sync.py for better tracking of synchronization processes. - Updated configuration mappings in config.py to include additional Access database tables. - Improved error handling and user feedback in vbareplace.py, including the ability to refresh linked tables in Access. - Adjusted polling intervals and batch sizes for performance optimization.
This commit is contained in:
@@ -6,182 +6,162 @@ import config
|
||||
import db_utils
|
||||
import logging
|
||||
from logging.handlers import TimedRotatingFileHandler
|
||||
from collections import defaultdict
|
||||
|
||||
# ================= 日志系统配置 (新增) =================
|
||||
|
||||
# ================= 日志系统配置 =================
|
||||
def setup_logger():
|
||||
"""
|
||||
配置双向日志系统:控制台输出 + 按天轮转的文件日志
|
||||
"""
|
||||
# 1. 确保日志目录存在
|
||||
log_dir = os.path.join(os.path.dirname(os.path.abspath(__file__)), 'log')
|
||||
if not os.path.exists(log_dir):
|
||||
os.makedirs(log_dir)
|
||||
|
||||
# 2. 创建 Logger
|
||||
logger = logging.getLogger("SyncService")
|
||||
logger.setLevel(logging.INFO)
|
||||
logger.handlers = [] # 清除旧句柄,防止重复打印
|
||||
logger.handlers = []
|
||||
|
||||
# 3. 文件处理器 (按天轮转)
|
||||
# filename: 正在写入的日志名
|
||||
# when='midnight': 每天午夜轮转
|
||||
# interval=1: 间隔1天
|
||||
# backupCount=30: 保留最近30个文件
|
||||
log_file_path = os.path.join(log_dir, 'sync_service.log')
|
||||
file_handler = TimedRotatingFileHandler(
|
||||
log_file_path, when='midnight', interval=1, backupCount=30, encoding='utf-8'
|
||||
)
|
||||
file_handler.suffix = "%Y-%m-%d" # 轮转后的文件名后缀
|
||||
# 文件日志格式:[时间] [级别] 消息
|
||||
file_handler.suffix = "%Y-%m-%d"
|
||||
file_fmt = logging.Formatter('%(asctime)s - [%(levelname)s] - %(message)s')
|
||||
file_handler.setFormatter(file_fmt)
|
||||
|
||||
# 4. 控制台处理器 (屏幕输出)
|
||||
console_handler = logging.StreamHandler(sys.stdout)
|
||||
# 控制台日志格式:[时间] 消息 (保持简洁)
|
||||
console_fmt = logging.Formatter('%(asctime)s - %(message)s', datefmt='%H:%M:%S')
|
||||
console_handler.setFormatter(console_fmt)
|
||||
|
||||
# 5. 添加处理器
|
||||
logger.addHandler(file_handler)
|
||||
logger.addHandler(console_handler)
|
||||
|
||||
return logger
|
||||
|
||||
# 初始化全局 Logger
|
||||
logger = setup_logger()
|
||||
|
||||
# ================= 辅助函数 =================
|
||||
|
||||
def clean_access_path(raw_addr):
|
||||
"""清洗 TableAddress 字段"""
|
||||
if not raw_addr: return ""
|
||||
s = str(raw_addr).strip()
|
||||
if s.upper().startswith(";DATABASE="): return s[10:]
|
||||
if s.upper().startswith("LOCAL:"): return s[6:]
|
||||
return s
|
||||
|
||||
def get_norm_path(path):
|
||||
"""获取标准化路径"""
|
||||
if not path: return ""
|
||||
return os.path.normpath(path).lower()
|
||||
|
||||
# ================= 主逻辑 =================
|
||||
|
||||
def process_logs():
|
||||
def process_sync_task():
|
||||
"""
|
||||
逻辑重构:
|
||||
遍历 config.SYNC_MAPPING 中的每一个文件 -> 去日志表查询该文件下特定表的未同步记录。
|
||||
"""
|
||||
sql_conn = db_utils.get_sql_conn()
|
||||
sql_cursor = sql_conn.cursor()
|
||||
sql_cursor.fast_executemany = True
|
||||
|
||||
cols = config.LOG_TABLE_CONFIG
|
||||
sql_cursor.fast_executemany = True # 高性能开关
|
||||
|
||||
# --- 0. 预处理配置映射 ---
|
||||
normalized_mapping = {}
|
||||
for cfg_path, cfg_tables in config.SYNC_MAPPING.items():
|
||||
norm_key = get_norm_path(cfg_path)
|
||||
normalized_mapping[norm_key] = {
|
||||
'original_path': cfg_path,
|
||||
'tables': cfg_tables
|
||||
}
|
||||
# 标记是否有工作被处理(用于控制轮询休眠时间)
|
||||
work_done = False
|
||||
|
||||
cols = config.LOG_TABLE_CONFIG
|
||||
log_full_name = db_utils.fmt_table(cols['schema'], cols['table_name'])
|
||||
|
||||
try:
|
||||
log_full_name = db_utils.fmt_table(cols['schema'], cols['table_name'])
|
||||
|
||||
# --- 1. 读取日志 ---
|
||||
query_log = f"""
|
||||
SELECT TOP 1000
|
||||
{cols['col_log_id']},
|
||||
{cols['col_table_name']},
|
||||
{cols['col_record_id']},
|
||||
{cols['col_address']}
|
||||
FROM {log_full_name}
|
||||
WHERE {cols['col_synced']} = 0
|
||||
AND TableType = 'LINKED_ACCESS'
|
||||
ORDER BY {cols['col_log_id']} ASC
|
||||
"""
|
||||
sql_cursor.execute(query_log)
|
||||
logs = sql_cursor.fetchall()
|
||||
|
||||
if not logs:
|
||||
return False
|
||||
|
||||
logger.info(f"📥 收到 {len(logs)} 条链接表变更日志...")
|
||||
|
||||
# --- 2. 任务分组 ---
|
||||
tasks = defaultdict(lambda: defaultdict(lambda: {'record_ids': set(), 'log_ids': []}))
|
||||
|
||||
for row in logs:
|
||||
log_id, acc_table, record_id, raw_address = row
|
||||
# === 核心循环:以配置文件为驱动 ===
|
||||
for clean_path, tables_map in config.SYNC_MAPPING.items():
|
||||
|
||||
clean_path = clean_access_path(raw_address)
|
||||
norm_lookup_key = get_norm_path(clean_path)
|
||||
|
||||
# 校验配置
|
||||
if norm_lookup_key not in normalized_mapping:
|
||||
# 记录一条警告日志,但不刷屏
|
||||
logger.warning(f"忽略未配置的文件路径: {clean_path}")
|
||||
continue
|
||||
|
||||
mapping_node = normalized_mapping[norm_lookup_key]
|
||||
|
||||
if acc_table not in mapping_node['tables']:
|
||||
logger.warning(f"忽略未配置的表: {acc_table} (在文件 {os.path.basename(clean_path)} 中)")
|
||||
# 1. 准备查询条件
|
||||
# 获取该文件下所有需要同步的表名列表
|
||||
target_tables = list(tables_map.keys())
|
||||
if not target_tables:
|
||||
continue
|
||||
|
||||
tasks[clean_path][acc_table]['record_ids'].add(record_id)
|
||||
tasks[clean_path][acc_table]['log_ids'].append(log_id)
|
||||
|
||||
# --- 3. 执行同步循环 ---
|
||||
for file_path, tables_data in tasks.items():
|
||||
# 构造 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 (兼容没有前缀的情况)
|
||||
]
|
||||
|
||||
if not os.path.exists(file_path):
|
||||
logger.error(f"无法访问文件: {file_path}")
|
||||
# 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 # 标记有工作
|
||||
logger.info(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):
|
||||
logger.error(f"无法访问文件: {clean_path}")
|
||||
continue
|
||||
|
||||
try:
|
||||
acc_conn = db_utils.get_access_conn(file_path)
|
||||
acc_conn = db_utils.get_access_conn(clean_path)
|
||||
acc_cursor = acc_conn.cursor()
|
||||
except Exception as conn_err:
|
||||
logger.error(f"连接 Access 失败 [{file_path}]: {conn_err}")
|
||||
logger.error(f"连接 Access 失败: {conn_err}")
|
||||
continue
|
||||
|
||||
norm_key = get_norm_path(file_path)
|
||||
config_node = normalized_mapping[norm_key]
|
||||
|
||||
for acc_table, data in tables_data.items():
|
||||
for acc_table, data in table_tasks.items():
|
||||
record_ids = list(data['record_ids'])
|
||||
log_ids = data['log_ids']
|
||||
|
||||
target_conf = config_node['tables'][acc_table]
|
||||
# 读取目标配置
|
||||
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)
|
||||
|
||||
try:
|
||||
# --- A. Access 查新数据 ---
|
||||
ids_placeholders = ','.join(['?'] * len(record_ids))
|
||||
|
||||
# A. Access 读取
|
||||
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 写入
|
||||
# --- 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)
|
||||
try: sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} ON")
|
||||
except: pass
|
||||
|
||||
sql_cursor.executemany(insert_sql, new_rows)
|
||||
|
||||
try: sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} OFF")
|
||||
except: pass
|
||||
|
||||
# C. 更新日志状态
|
||||
# 3. 标记日志 Synced = 1
|
||||
log_placeholders = ','.join(['?'] * len(log_ids))
|
||||
update_log_sql = f"""
|
||||
UPDATE {log_full_name}
|
||||
@@ -191,15 +171,15 @@ def process_logs():
|
||||
sql_cursor.execute(update_log_sql, log_ids)
|
||||
|
||||
sql_conn.commit()
|
||||
logger.info(f"✅ 同步成功: {target_table} 更新 {len(record_ids)} 条 (来源: {os.path.basename(file_path)})")
|
||||
logger.info(f" -> {target_table} 同步成功: {len(record_ids)} 条")
|
||||
|
||||
except Exception as tbl_err:
|
||||
logger.error(f"处理表 {acc_table} 异常: {tbl_err}")
|
||||
logger.error(f" -> 处理表 {acc_table} 异常: {tbl_err}")
|
||||
sql_conn.rollback()
|
||||
|
||||
acc_conn.close()
|
||||
acc_conn.close() # 关闭 Access 连接
|
||||
|
||||
return True
|
||||
return work_done
|
||||
|
||||
except Exception as e:
|
||||
logger.critical(f"全局严重异常: {e}")
|
||||
@@ -209,14 +189,18 @@ def process_logs():
|
||||
except: pass
|
||||
|
||||
if __name__ == "__main__":
|
||||
logger.info("=== 增量同步服务已启动 (日志模式: log/sync_service.log) ===")
|
||||
logger.info("=== 增量同步服务启动 (配置驱动模式) ===")
|
||||
logger.info("正在按 config.SYNC_MAPPING 轮询...")
|
||||
|
||||
while True:
|
||||
try:
|
||||
has_work = process_logs()
|
||||
time.sleep(0.5 if has_work else config.POLL_INTERVAL)
|
||||
has_work = process_sync_task()
|
||||
# 如果有工作,说明可能还有积压,休息短一点(0.1s)
|
||||
# 如果没工作,休息标准间隔(5s)
|
||||
time.sleep(0.1 if has_work else config.POLL_INTERVAL)
|
||||
except KeyboardInterrupt:
|
||||
logger.info("收到停止指令,服务正在关闭...")
|
||||
logger.info("服务停止")
|
||||
break
|
||||
except Exception as e:
|
||||
logger.critical(f"主循环崩溃: {e}")
|
||||
time.sleep(5) # 防止死循环刷屏
|
||||
time.sleep(5)
|
||||
Reference in New Issue
Block a user