Refactor logging functions to enhance clarity and consistency in log messages
This commit is contained in:
@@ -35,6 +35,55 @@ def setup_logger():
|
||||
|
||||
logger = setup_logger()
|
||||
|
||||
# ================= 日志辅助函数 =================
|
||||
def log_success(message):
|
||||
"""成功消息 - 绿色感觉"""
|
||||
logger.info(f"✅ [成功] {message}")
|
||||
|
||||
def log_error(message):
|
||||
"""错误消息 - 红色感觉"""
|
||||
logger.error(f"❌ [错误] {message}")
|
||||
|
||||
def log_warning(message):
|
||||
"""警告消息 - 黄色感觉"""
|
||||
logger.warning(f"⚠️ [警告] {message}")
|
||||
|
||||
def log_info(message):
|
||||
"""信息消息"""
|
||||
logger.info(f"ℹ️ {message}")
|
||||
|
||||
def log_processing(message):
|
||||
"""处理中消息"""
|
||||
logger.info(f"🔄 [处理] {message}")
|
||||
|
||||
def log_skip(message):
|
||||
"""跳过消息"""
|
||||
logger.info(f"⏭️ [跳过] {message}")
|
||||
|
||||
def log_critical(message):
|
||||
"""严重错误"""
|
||||
logger.critical(f"🔥 [严重] {message}")
|
||||
|
||||
def log_start(message):
|
||||
"""启动消息"""
|
||||
logger.info(f"🚀 [启动] {message}")
|
||||
|
||||
def log_stop(message):
|
||||
"""停止消息"""
|
||||
logger.info(f"🛑 [停止] {message}")
|
||||
|
||||
def log_file(message):
|
||||
"""文件操作消息"""
|
||||
logger.info(f"📂 [文件] {message}")
|
||||
|
||||
def log_database(message):
|
||||
"""数据库操作消息"""
|
||||
logger.info(f"💾 [数据库] {message}")
|
||||
|
||||
def log_sync(message):
|
||||
"""同步操作消息"""
|
||||
logger.info(f"🔃 [同步] {message}")
|
||||
|
||||
# ================= 主逻辑 =================
|
||||
|
||||
def process_sync_task():
|
||||
@@ -102,7 +151,7 @@ def process_sync_task():
|
||||
continue # 这个文件没有需要同步的记录,检查下一个文件
|
||||
|
||||
work_done = True # 标记有工作
|
||||
logger.info(f"📂 文件 {os.path.basename(clean_path)} 发现 {len(logs)} 条变更...")
|
||||
log_file(f"{os.path.basename(clean_path)} 发现 {len(logs)} 条待同步变更")
|
||||
|
||||
# 3. 本地分组 (按表名)
|
||||
# 结构: table_tasks[TableName] = { ids: {}, log_ids: [] }
|
||||
@@ -116,16 +165,22 @@ def process_sync_task():
|
||||
|
||||
# 4. 执行同步 (连接一次 Access,处理多张表)
|
||||
if not os.path.exists(clean_path):
|
||||
logger.error(f"无法访问文件: {clean_path}")
|
||||
log_error(f"无法访问文件: {clean_path}")
|
||||
continue
|
||||
|
||||
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:
|
||||
logger.error(f"连接 Access 失败: {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']
|
||||
@@ -137,6 +192,8 @@ def process_sync_task():
|
||||
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)} 条记录)")
|
||||
|
||||
try:
|
||||
# --- A. Access 查新数据 ---
|
||||
ids_placeholders = ','.join(['?'] * len(record_ids))
|
||||
@@ -171,26 +228,40 @@ def process_sync_task():
|
||||
sql_cursor.execute(update_log_sql, log_ids)
|
||||
|
||||
sql_conn.commit()
|
||||
logger.info(f" -> {target_table} 同步成功: {len(record_ids)} 条")
|
||||
log_success(f"表 [{target_table}] 同步完成: {len(record_ids)} 条记录")
|
||||
|
||||
file_success_count += 1
|
||||
file_total_records += len(record_ids)
|
||||
|
||||
except Exception as tbl_err:
|
||||
logger.error(f" -> 处理表 {acc_table} 异常: {tbl_err}")
|
||||
log_error(f"表 [{acc_table}] 同步失败: {tbl_err}")
|
||||
sql_conn.rollback()
|
||||
file_error_count += 1
|
||||
|
||||
acc_conn.close() # 关闭 Access 连接
|
||||
|
||||
# 输出文件级别的汇总
|
||||
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:
|
||||
logger.critical(f"全局严重异常: {e}")
|
||||
log_critical(f"全局异常: {e}")
|
||||
return False
|
||||
finally:
|
||||
try: sql_conn.close()
|
||||
except: pass
|
||||
|
||||
if __name__ == "__main__":
|
||||
logger.info("=== 增量同步服务启动 (配置驱动模式) ===")
|
||||
logger.info("正在按 config.SYNC_MAPPING 轮询...")
|
||||
log_start("增量同步服务已启动 (配置驱动模式)")
|
||||
log_info(f"轮询间隔: {config.POLL_INTERVAL} 秒")
|
||||
log_info(f"监控配置: {len(config.SYNC_MAPPING)} 个文件")
|
||||
logger.info("=" * 70)
|
||||
|
||||
while True:
|
||||
try:
|
||||
@@ -199,8 +270,9 @@ if __name__ == "__main__":
|
||||
# 如果没工作,休息标准间隔(5s)
|
||||
time.sleep(0.1 if has_work else config.POLL_INTERVAL)
|
||||
except KeyboardInterrupt:
|
||||
logger.info("服务停止")
|
||||
logger.info("=" * 70)
|
||||
log_stop("收到停止信号,服务正在关闭...")
|
||||
break
|
||||
except Exception as e:
|
||||
logger.critical(f"主循环崩溃: {e}")
|
||||
log_critical(f"主循环崩溃: {e}")
|
||||
time.sleep(5)
|
||||
Reference in New Issue
Block a user