diff --git a/.claude/settings.local.json b/.claude/settings.local.json index 2a56823..bb10ce6 100644 --- a/.claude/settings.local.json +++ b/.claude/settings.local.json @@ -7,7 +7,8 @@ "Bash(python -c \"import etl_manager; print\\(''etl_manager.py: import OK''\\)\")", "Bash(python -c \"import init_full_sync; print\\(''init_full_sync.py: import OK''\\)\")", "Bash(python -c \"import migration; print\\(''migration.py: import OK''\\)\")", - "Bash(python migration.py)" + "Bash(python migration.py)", + "Bash(python -c \"import log_utils; from log_utils import LoggerManager, log_start, log_complete; print\\(''log_utils: OK''\\)\")" ] } } diff --git a/config/database.py b/config/database.py index e24e4d8..e41bd7f 100644 --- a/config/database.py +++ b/config/database.py @@ -8,7 +8,7 @@ SQL_SERVER_CONFIG = { 'database': 'CompanyDB', 'username': 'peng', 'password': 'Cqbld123456.', - 'trust_server_certificate': 'yes' + 'TrustServerCertificate': 'yes' } # ================= Access 驱动配置 ================= diff --git a/etl_manager.py b/etl_manager.py index 6ac3493..0932296 100644 --- a/etl_manager.py +++ b/etl_manager.py @@ -13,26 +13,13 @@ from sqlalchemy.engine import URL from sqlalchemy.types import NVARCHAR, Integer, Date # 导入配置 +from log_utils import (log_error, log_warning, log_info, log_processing, log_file, log_sync, + log_start, log_complete, LoggerManager) from config import DB_CONFIG, CACHE_DIR, EXCEL_CONFIGS, BATCH_SIZE, TABLE_SCHEMA # ================= 抑制 openpyxl 的数据验证警告 ================= warnings.filterwarnings('ignore', category=UserWarning, module='openpyxl') -# ================= 日志配置 ================= -# 配置控制台输出使用 UTF-8 编码,确保中文正确显示 -console_handler = logging.StreamHandler(sys.stdout) -console_handler.setFormatter(logging.Formatter('%(asctime)s - %(levelname)s - %(message)s')) - -logging.basicConfig( - level=logging.INFO, - format='%(asctime)s - %(levelname)s - %(message)s', - handlers=[ - console_handler, - logging.FileHandler("sync_log.txt", encoding='utf-8') - ] -) -logger = logging.getLogger(__name__) - class DataSynchronizer: def __init__(self, force_sync=False): self.force_sync = force_sync @@ -49,7 +36,7 @@ class DataSynchronizer: f"DATABASE={DB_CONFIG['database']};" f"UID={DB_CONFIG['username']};" f"PWD={DB_CONFIG['password']};" - f"TrustServerCertificate={DB_CONFIG.get('TrustServerCertificate', 'no')};" + f"TrustServerCertificate={DB_CONFIG.get('TrustServerCertificate', 'yes')};" ) connection_url = URL.create("mssql+pyodbc", query={"odbc_connect": connection_string}) return create_engine(connection_url, fast_executemany=True) @@ -67,7 +54,7 @@ class DataSynchronizer: if remote_mtime > local_mtime + 1: return True, f"源文件更新" except OSError as e: - logger.error(f"无法访问源文件: {remote_path}, Error: {e}") + log_error(f"无法访问源文件: {remote_path}, Error: {e}") return False, "源文件无法访问" return False, "文件未变更" @@ -82,7 +69,7 @@ class DataSynchronizer: df = df.dropna(subset=['总排号']) df = df[df['总排号'].astype(str).str.strip() != ''] else: - logger.error("数据源中找不到映射后的[总排号]列,跳过此 sheet") + log_error("数据源中找不到映射后的[总排号]列,跳过此 sheet") return None, None # ★ 新增:去除重复的总排号(保留第一条) @@ -90,7 +77,7 @@ class DataSynchronizer: df['总排号'] = df['总排号'].astype(str).str.strip() duplicates = df[df.duplicated(subset=['总排号'], keep='first')] if not duplicates.empty: - logger.warning(f"发现 {len(duplicates)} 条重复的总排号,已自动去重。重复的总排号: {duplicates['总排号'].tolist()[:10]}") + log_warning(f"发现 {len(duplicates)} 条重复的总排号,已自动去重。重复的总排号: {duplicates['总排号'].tolist()[:10]}") df = df.drop_duplicates(subset=['总排号'], keep='first') # 3. 补全列 @@ -156,19 +143,19 @@ class DataSynchronizer: df_update = df[df['总排号'].isin(existing_id_set)].copy() df_insert = df[~df['总排号'].isin(existing_id_set)].copy() - logger.info(f"分析结果: 需插入 {len(df_insert)} 条, 需更新 {len(df_update)} 条") + log_info(f"分析结果: 需插入 {len(df_insert)} 条, 需更新 {len(df_update)} 条") # 1. 插入新数据 if not df_insert.empty: - logger.info("正在执行批量插入...") + log_info("正在执行批量插入...") df_insert.to_sql('executionCardData', self.engine, schema='warehouseOutbound', if_exists='append', index=False, chunksize=BATCH_SIZE, dtype=dtype_dict) - logger.info("批量插入完成。") + log_info("批量插入完成。") # 2. 更新现有数据 - 改用逐条或小批量 UPDATE if not df_update.empty: - logger.info("正在执行批量更新...") + log_info("正在执行批量更新...") cols = [c for c in df.columns if c != '总排号'] set_clause = ", ".join([f"[{c}] = :{c}" for c in cols]) @@ -191,11 +178,11 @@ class DataSynchronizer: update_count += result.rowcount if (i + batch_size) % 5000 == 0: - logger.info(f"已更新 {i + batch_size}/{total_rows} 条记录...") + log_info(f"已更新 {i + batch_size}/{total_rows} 条记录...") # 修复 SQL Server executemany 返回负数 rowcount 的问题 affected_rows = abs(update_count) if update_count < 0 else total_rows - logger.info(f"批量更新完成,共影响 {affected_rows} 行。") + log_info(f"批量更新完成,共影响 {affected_rows} 行。") def process_excel_files(self): for cfg in EXCEL_CONFIGS: @@ -206,7 +193,7 @@ class DataSynchronizer: should_sync, reason = self._should_process_file(remote_path, local_path) if should_sync: - logger.info(f"开始处理文件: {filename} ({reason})") + log_info(f"开始处理文件: {filename} ({reason})") try: # 复制文件到本地缓存(只复制一次) if os.path.exists(remote_path): @@ -214,7 +201,7 @@ class DataSynchronizer: # 遍历该文件的所有指定 sheet for sheet_name in cfg['sheet_names']: - logger.info(f" → 处理工作表: {sheet_name} (合同年份: {cfg['contract_year']})") + log_info(f" → 处理工作表: {sheet_name} (合同年份: {cfg['contract_year']})") try: df = pd.read_excel(local_path, sheet_name=sheet_name, header=0, engine='openpyxl') df.columns = [str(c).strip() for c in df.columns] @@ -224,21 +211,21 @@ class DataSynchronizer: if cleaned_df is not None: self._sync_to_db(cleaned_df, dtype_mapping) - logger.info(f" 工作表 {sheet_name} 同步成功。") + log_info(f" 工作表 {sheet_name} 同步成功。") else: - logger.warning(f" 工作表 {sheet_name} 清洗失败,跳过。") + log_warning(f" 工作表 {sheet_name} 清洗失败,跳过。") except Exception as e: - logger.error(f" 处理工作表 {sheet_name} 时发生错误: {str(e)}", exc_info=True) + log_error(f" 处理工作表 {sheet_name} 时发生错误: {str(e)}", exc_info=True) - logger.info(f"文件 {filename} 所有工作表处理完成。") + log_info(f"文件 {filename} 所有工作表处理完成。") except Exception as e: - logger.error(f"处理文件 {filename} 时发生错误: {str(e)}", exc_info=True) + log_error(f"处理文件 {filename} 时发生错误: {str(e)}", exc_info=True) else: - logger.info(f"跳过文件: {filename} ({reason})") + log_info(f"跳过文件: {filename} ({reason})") def generate_contract_data(self): - logger.info("开始生成/更新 contractData 表...") + log_info("开始生成/更新 contractData 表...") merge_sql = """ WITH SourceData AS ( @@ -299,21 +286,30 @@ class DataSynchronizer: try: with self.engine.begin() as conn: result = conn.execute(text(merge_sql)) - logger.info(f"ContractData 表同步完成 (SQL Server 内部处理)。rowcount: {result.rowcount}") + log_info(f"ContractData 表同步完成 (SQL Server 内部处理)。rowcount: {result.rowcount}") except Exception as e: - logger.error(f"生成 ContractData 失败: {e}", exc_info=True) + log_error(f"生成 ContractData 失败: {e}", exc_info=True) def main(): + # 初始化日志管理器 + LoggerManager("etl_manager", log_prefix="sync") + parser = argparse.ArgumentParser(description="Excel数据同步至SQL Server") parser.add_argument('--force', action='store_true', help='强制同步所有文件') args = parser.parse_args() syncer = DataSynchronizer(force_sync=args.force) - logger.info("================= 任务开始 =================") + + if args.force: + log_start("Excel 同步任务 (强制模式)") + else: + log_start("Excel 同步任务 (增量模式)") + syncer.process_excel_files() syncer.generate_contract_data() - logger.info("================= 任务结束 =================") + + log_complete("Excel 同步任务已完成") if __name__ == "__main__": main() \ No newline at end of file diff --git a/init_full_sync.py b/init_full_sync.py index 0e6c334..3eb16c8 100644 --- a/init_full_sync.py +++ b/init_full_sync.py @@ -3,79 +3,11 @@ from config import SQL_SERVER_CONN, ACCESS_DRIVER, SYNC_MAPPING, BATCH_SIZE import db_utils import time import os -import logging -from datetime import datetime +from log_utils import (LoggerManager, log_success, log_error, log_warning, log_info, log_processing, + log_skip, log_start, log_file, log_database, log_sync, log_complete) -def setup_logger(): - """初始化日志配置,每次运行生成独立的日志文件""" - # 创建log目录 - log_dir = "log" - if not os.path.exists(log_dir): - os.makedirs(log_dir) - - # 生成带时间戳的日志文件名 - timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") - log_file = os.path.join(log_dir, f"full_sync_{timestamp}.log") - - # 配置日志格式 - logging.basicConfig( - level=logging.INFO, - format='%(asctime)s [%(levelname)s] %(message)s', - datefmt='%Y-%m-%d %H:%M:%S', - handlers=[ - logging.FileHandler(log_file, encoding='utf-8'), - logging.StreamHandler() # 同时输出到控制台 - ] - ) - - logger = logging.getLogger(__name__) - logger.info(f"📝 日志文件: {log_file}") - return logger - -# ================= 日志辅助函数 ================= -def log_success(logger, message): - """成功消息""" - logger.info(f"✅ [成功] {message}") - -def log_error(logger, message, exc_info=False): - """错误消息""" - logger.error(f"❌ [错误] {message}", exc_info=exc_info) - -def log_warning(logger, message): - """警告消息""" - logger.warning(f"⚠️ [警告] {message}") - -def log_info(logger, message): - """信息消息""" - logger.info(f"ℹ️ {message}") - -def log_processing(logger, message): - """处理中消息""" - logger.info(f"🔄 [处理] {message}") - -def log_skip(logger, message): - """跳过消息""" - logger.info(f"⏭️ [跳过] {message}") - -def log_start(logger, message): - """启动消息""" - logger.info(f"🚀 [启动] {message}") - -def log_file(logger, message): - """文件操作消息""" - logger.info(f"📂 [文件] {message}") - -def log_database(logger, message): - """数据库操作消息""" - logger.info(f"💾 [数据库] {message}") - -def log_sync(logger, message): - """同步操作消息""" - logger.info(f"🔃 [同步] {message}") - -def log_complete(logger, message): - """完成消息""" - logger.info(f"🎯 [完成] {message}") +# 初始化日志管理器 +LoggerManager("init_full_sync", log_prefix="full_sync") def format_duration(seconds): """格式化时长显示""" @@ -99,13 +31,11 @@ def has_identity_column(sql_cursor, schema, table): return sql_cursor.fetchone()[0] > 0 def run_full_sync(): - logger = setup_logger() - - logger.info("=" * 70) - log_start(logger, "全量初始化同步") - log_info(logger, f"配置文件数: {len(SYNC_MAPPING)} 个") - log_info(logger, f"批次大小: {BATCH_SIZE} 行") - logger.info("=" * 70) + log_info("=" * 70) + log_start( "全量初始化同步") + log_info( f"配置文件数: {len(SYNC_MAPPING)} 个") + log_info( f"批次大小: {BATCH_SIZE} 行") + log_info("=" * 70) sync_start_time = time.time() total_tables = 0 @@ -117,18 +47,18 @@ def run_full_sync(): sql_conn = db_utils.get_sql_conn() sql_cursor = sql_conn.cursor() sql_cursor.fast_executemany = True - log_database(logger, "SQL Server 连接成功") + log_database( "SQL Server 连接成功") except Exception as e: - log_error(logger, f"SQL Server 连接失败: {e}") + log_error( f"SQL Server 连接失败: {e}") return # 第一层循环:遍历配置文件中的所有文件路径 for acc_path, tables_map in SYNC_MAPPING.items(): if not os.path.exists(acc_path): - log_skip(logger, f"文件不存在: {acc_path}") + log_skip( f"文件不存在: {acc_path}") continue - log_file(logger, f"开始处理: {os.path.basename(acc_path)}") + log_file( f"开始处理: {os.path.basename(acc_path)}") file_start_time = time.time() file_success = 0 file_failed = 0 @@ -137,7 +67,7 @@ def run_full_sync(): try: acc_conn = db_utils.get_access_conn(acc_path) acc_cursor = acc_conn.cursor() - log_database(logger, f"Access 文件连接成功: {os.path.basename(acc_path)}") + log_database( f"Access 文件连接成功: {os.path.basename(acc_path)}") # 第二层循环:遍历该文件下的所有表映射 for acc_table, target_config in tables_map.items(): @@ -146,7 +76,7 @@ def run_full_sync(): target_table = target_config['target_table'] full_target_name = db_utils.fmt_table(target_schema, target_table) - log_processing(logger, f"表 [{acc_table}] → [{target_table}]") + log_processing( f"表 [{acc_table}] → [{target_table}]") table_start_time = time.time() try: @@ -155,12 +85,12 @@ def run_full_sync(): # 2. 获取列结构 columns = [col[0] for col in acc_cursor.description] - log_info(logger, f" 检测到 {len(columns)} 个列") + log_info( f" 检测到 {len(columns)} 个列") insert_sql = db_utils.generate_insert_sql(target_schema, target_table, columns) # 3. TRUNCATE 目标表 sql_cursor.execute(f"TRUNCATE TABLE {full_target_name}") - log_info(logger, f" 已清空目标表") + log_info( f" 已清空目标表") # 4. 检查并启用 IDENTITY_INSERT has_identity = has_identity_column(sql_cursor, target_schema, target_table) @@ -170,13 +100,13 @@ def run_full_sync(): try: sql_cursor.execute(f"SET IDENTITY_INSERT {full_target_name} ON") identity_enabled = True - log_info(logger, f" 已启用 IDENTITY_INSERT") + log_info( f" 已启用 IDENTITY_INSERT") except Exception as id_err: - log_error(logger, f" 无法启用 IDENTITY_INSERT: {id_err}") + log_error( f" 无法启用 IDENTITY_INSERT: {id_err}") raise # 5. 传输数据 - log_info(logger, f" 开始数据传输...") + log_info( f" 开始数据传输...") acc_cursor.execute(f"SELECT * FROM [{acc_table}]") table_rows = 0 @@ -195,23 +125,23 @@ def run_full_sync(): if (current_time - last_log_time >= 5) or (table_rows % 10000 == 0): elapsed = current_time - start_time rate = table_rows / elapsed if elapsed > 0 else 0 - log_info(logger, f" 进度: {table_rows:,} 行 | 速率: {rate:,.0f} 行/秒") + log_info( f" 进度: {table_rows:,} 行 | 速率: {rate:,.0f} 行/秒") last_log_time = current_time # 6. 关闭 IDENTITY_INSERT if identity_enabled: try: sql_cursor.execute(f"SET IDENTITY_INSERT {full_target_name} OFF") - log_info(logger, f" 已关闭 IDENTITY_INSERT") + log_info( f" 已关闭 IDENTITY_INSERT") except Exception as id_err: - log_warning(logger, f" 关闭 IDENTITY_INSERT 时警告: {id_err}") + log_warning( f" 关闭 IDENTITY_INSERT 时警告: {id_err}") sql_conn.commit() # 显示最终统计 table_time = time.time() - table_start_time final_rate = table_rows / table_time if table_time > 0 else 0 - log_success(logger, f"表 [{target_table}] 完成: {table_rows:,} 行 | 速率: {final_rate:,.0f} 行/秒 | 用时: {format_duration(table_time)}") + log_success( f"表 [{target_table}] 完成: {table_rows:,} 行 | 速率: {final_rate:,.0f} 行/秒 | 用时: {format_duration(table_time)}") success_tables += 1 file_success += 1 @@ -221,7 +151,7 @@ def run_full_sync(): except Exception as tbl_err: failed_tables += 1 file_failed += 1 - log_error(logger, f"表 [{acc_table}] 同步失败: {tbl_err}", exc_info=True) + log_error( f"表 [{acc_table}] 同步失败: {tbl_err}", exc_info=True) sql_conn.rollback() # 确保清理 IDENTITY_INSERT 状态 @@ -239,24 +169,24 @@ def run_full_sync(): if file_failed > 0: summary += f" | 失败 {file_failed} 张表" summary += f" | 共 {file_rows:,} 行 | 用时: {format_duration(file_time)}" - log_sync(logger, summary) + log_sync( summary) except Exception as file_err: - log_error(logger, f"文件 [{os.path.basename(acc_path)}] 处理失败: {file_err}", exc_info=True) + log_error( f"文件 [{os.path.basename(acc_path)}] 处理失败: {file_err}", exc_info=True) sql_conn.close() # 总结统计 sync_time = time.time() - sync_start_time - logger.info("\n" + "=" * 70) - log_complete(logger, "全量同步任务结束") - logger.info("-" * 70) - log_info(logger, f"总表数: {total_tables} 张 | 成功: {success_tables} 张 | 失败: {failed_tables} 张") - log_info(logger, f"总行数: {total_rows:,} 行") - log_info(logger, f"总用时: {format_duration(sync_time)}") + log_info("\n" + "=" * 70) + log_complete( "全量同步任务结束") + log_info("-" * 70) + log_info( f"总表数: {total_tables} 张 | 成功: {success_tables} 张 | 失败: {failed_tables} 张") + log_info( f"总行数: {total_rows:,} 行") + log_info( f"总用时: {format_duration(sync_time)}") if total_rows > 0 and sync_time > 0: - log_info(logger, f"整体平均速率: {total_rows/sync_time:,.0f} 行/秒") - logger.info("=" * 70) + log_info( f"整体平均速率: {total_rows/sync_time:,.0f} 行/秒") + log_info("=" * 70) if __name__ == "__main__": run_full_sync() \ No newline at end of file diff --git a/log_utils.py b/log_utils.py new file mode 100644 index 0000000..e704cda --- /dev/null +++ b/log_utils.py @@ -0,0 +1,136 @@ +# log_utils.py +# 统一的日志工具模块 + +import logging +import os +import sys +from datetime import datetime +import ntfy_utils + +# ================= 全局 logger 实例 ================= +_logger = None + +# ================= 日志格式常量 ================= +LOG_FORMAT = '%(asctime)s [%(levelname)s] %(message)s' +DATE_FORMAT = '%Y-%m-%d %H:%M:%S' + + +class LoggerManager: + """统一日志管理器""" + + def __init__(self, name, log_prefix="app", log_dir="log"): + """ + 初始化日志管理器 + + Args: + name: logger 名称 + log_prefix: 日志文件前缀(如 app, sync, migration) + log_dir: 日志目录 + """ + global _logger + + # 创建日志目录 + log_path = os.path.join(os.getcwd(), log_dir) + os.makedirs(log_path, exist_ok=True) + + # 创建带时间戳的日志文件 + timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") + log_file = os.path.join(log_path, f"{log_prefix}_{timestamp}.log") + + # 创建 logger + _logger = logging.getLogger(name) + _logger.setLevel(logging.INFO) + _logger.handlers = [] + + # 文件处理器 + file_handler = logging.FileHandler(log_file, encoding='utf-8') + file_formatter = logging.Formatter(LOG_FORMAT, DATE_FORMAT) + file_handler.setFormatter(file_formatter) + _logger.addHandler(file_handler) + + # 控制台处理器 + console_handler = logging.StreamHandler(sys.stdout) + console_formatter = logging.Formatter(LOG_FORMAT, DATE_FORMAT) + console_handler.setFormatter(console_formatter) + _logger.addHandler(console_handler) + + _logger.info(f"日志文件: {log_file}") + + @staticmethod + def get_logger(): + """获取全局 logger 实例""" + return _logger + + +# ================= 统一的日志辅助函数 ================= +# 无需 logger 参数,内部使用全局 _logger + +def log_success(message): + """成功消息 - ✅ 发送 ntfy 通知 (default)""" + _logger.info(f"✅ [成功] {message}") + ntfy_utils.send_ntfy(f"✅ [成功] {message}", title="同步任务成功", priority="default", tags=["white_check_mark"]) + + +def log_error(message, exc_info=False): + """错误消息 - ❌ 发送 ntfy 通知 (high)""" + _logger.error(f"❌ [错误] {message}", exc_info=exc_info) + ntfy_utils.send_error(f"❌ [错误] {message}") + + +def log_warning(message): + """警告消息 - ⚠️ 不发送 ntfy 通知""" + _logger.warning(f"⚠️ [警告] {message}") + + +def log_critical(message, exc_info=False): + """严重错误 - 🔥 发送 ntfy 通知 (urgent)""" + _logger.critical(f"🔥 [严重] {message}", exc_info=exc_info) + ntfy_utils.send_critical(f"🔥 [严重] {message}") + + +def log_start(message): + """启动消息 - 🚀 发送 ntfy 通知 (default)""" + _logger.info(f"🚀 [启动] {message}") + ntfy_utils.send_ntfy(f"🚀 [启动] {message}", title="任务启动", priority="default", tags=["rocket"]) + + +def log_complete(message): + """完成消息 - 🎯 发送 ntfy 通知 (default)""" + _logger.info(f"🎯 [完成] {message}") + ntfy_utils.send_ntfy(f"🎯 [完成] {message}", title="任务完成", priority="default", tags=["checkered_flag"]) + + +def log_stop(message): + """停止消息 - 🛑 发送 ntfy 通知 (default)""" + _logger.info(f"🛑 [停止] {message}") + ntfy_utils.send_ntfy(f"🛑 [停止] {message}", title="服务停止", priority="default", tags=["stop_sign"]) + + +def log_info(message): + """信息消息 - ℹ️ 不发送 ntfy 通知""" + _logger.info(f"ℹ️ {message}") + + +def log_processing(message): + """处理中消息 - 🔄 不发送 ntfy 通知""" + _logger.info(f"🔄 [处理] {message}") + + +def log_skip(message): + """跳过消息 - ⏭️ 不发送 ntfy 通知""" + _logger.info(f"⏭️ [跳过] {message}") + + +def log_file(message): + """文件操作消息 - 📂 不发送 ntfy 通知""" + _logger.info(f"📂 [文件] {message}") + + +def log_database(message): + """数据库操作消息 - 💾 不发送 ntfy 通知""" + _logger.info(f"💾 [数据库] {message}") + + +def log_sync(message): + """同步操作消息 - 🔃 不发送 ntfy 通知""" + _logger.info(f"🔃 [同步] {message}") diff --git a/migration.py b/migration.py index 66def12..51f92fd 100644 --- a/migration.py +++ b/migration.py @@ -4,7 +4,10 @@ import shutil import urllib from sqlalchemy import create_engine, text from config import DB_CONFIG, MIGRATION_TASKS, TEMP_DIR -import ntfy_utils # 确保该文件在同一目录下 +from log_utils import LoggerManager, log_start, log_skip, log_processing, log_success, log_warning, log_error, log_complete + +# 初始化日志管理器 +LoggerManager("migration", log_prefix="migration") # ========================================== # 1. 脚本配置 (Configuration) @@ -67,7 +70,7 @@ def run_migration(): sync_count = 0 error_count = 0 - print(f"🚀 开始增量同步任务 (强制更新={FORCE_UPDATE})") + log_start(f"增量同步任务 (强制更新={FORCE_UPDATE})") for task in MIGRATION_TASKS: remote_path = task['file_path'] @@ -77,8 +80,7 @@ def run_migration(): # 1. 检查源文件 if not os.path.exists(remote_path): msg = f"远程文件未找到: {remote_path}" - print(f"❌ {msg}") - ntfy_utils.send_error(msg) + log_error(msg) continue # 2. 增量判定 @@ -86,10 +88,10 @@ def run_migration(): local_mtime = get_file_mtime(local_path) if not FORCE_UPDATE and os.path.exists(local_path) and remote_mtime <= local_mtime: - print(f"⏭️ 跳过: {filename} (文件未变更)") + log_skip(f"{filename} (文件未变更)") continue - print(f"🔄 正在处理: {filename} ...") + log_processing(f"正在处理: {filename} ...") try: # 3. 复制文件到本地 temp @@ -112,7 +114,7 @@ def run_migration(): missing = [c for c in source_cols if c not in df.columns] if missing: - print(f" ⚠️ Sheet[{sheet_name}] 缺失列: {missing}") + log_warning(f"Sheet[{sheet_name}] 缺失列: {missing}") continue # 提取并重命名 @@ -151,23 +153,19 @@ def run_migration(): chunksize=1000 ) - print(f" ✅ 成功同步: {len(final_df)} 行记录") + log_success(f"成功同步: {len(final_df)} 行记录") sync_count += 1 else: - print(f" ⚠️ 警告: 文件内容为空或格式不符") + log_warning("文件内容为空或格式不符") except Exception as e: error_msg = f"文件 [{filename}] 处理失败: {str(e)}" - print(f"❌ {error_msg}") - ntfy_utils.send_error(error_msg) + log_error(error_msg) error_count += 1 # 结束汇总 summary = f"同步完成: 成功 {sync_count} 个文件, 失败 {error_count} 个文件。" - print(f"\n🏁 {summary}") - if sync_count > 0: - # 只有在有实际更新时才发送成功通知 - ntfy_utils.send_ntfy(summary, title="📊 数据迁移报告", tags=["package"]) + log_complete(f"同步完成: 成功 {sync_count} 个文件, 失败 {error_count} 个文件") if __name__ == "__main__": run_migration() \ No newline at end of file diff --git a/run_incremental_sync.py b/run_incremental_sync.py index 0c3405e..29c14ed 100644 --- a/run_incremental_sync.py +++ b/run_incremental_sync.py @@ -4,91 +4,11 @@ import sys import pyodbc from config import SQL_SERVER_CONN, ACCESS_DRIVER, SYNC_MAPPING, POLL_INTERVAL, LOG_TABLE_CONFIG import db_utils -import logging -from logging.handlers import TimedRotatingFileHandler -import ntfy_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) -# ================= 日志系统配置 ================= -def setup_logger(): - log_dir = os.path.join(os.path.dirname(os.path.abspath(__file__)), 'log') - if not os.path.exists(log_dir): - os.makedirs(log_dir) - - logger = logging.getLogger("SyncService") - logger.setLevel(logging.INFO) - logger.handlers = [] - - 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_fmt = logging.Formatter('%(asctime)s - [%(levelname)s] - %(message)s') - file_handler.setFormatter(file_fmt) - - console_handler = logging.StreamHandler(sys.stdout) - console_fmt = logging.Formatter('%(asctime)s - %(message)s', datefmt='%H:%M:%S') - console_handler.setFormatter(console_fmt) - - logger.addHandler(file_handler) - logger.addHandler(console_handler) - return logger - -logger = setup_logger() - -# ================= 日志辅助函数 ================= -def log_success(message): - """成功消息 - 绿色感觉""" - logger.info(f"✅ [成功] {message}") - -def log_error(message): - """错误消息 - 红色感觉""" - logger.error(f"❌ [错误] {message}") - ntfy_utils.send_error(f"❌ [错误] {message}") - -def log_warning(message): - """警告消息 - 黄色感觉""" - logger.warning(f"⚠️ [警告] {message}") - ntfy_utils.send_error(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}") - ntfy_utils.send_critical(f"🔥 [严重] {message}") - -def log_start(message): - """启动消息""" - logger.info(f"🚀 [启动] {message}") - ntfy_utils.send_ntfy(f"🚀 [启动] {message}") - -def log_stop(message): - """停止消息""" - logger.info(f"🛑 [停止] {message}") - ntfy_utils.send_ntfy(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}") +# 初始化日志管理器 +LoggerManager("run_incremental_sync", log_prefix="incremental") # ================= 主逻辑 ================= @@ -267,7 +187,7 @@ if __name__ == "__main__": log_start("增量同步服务已启动 (配置驱动模式)") log_info(f"轮询间隔: {POLL_INTERVAL} 秒") log_info(f"监控配置: {len(SYNC_MAPPING)} 个文件") - logger.info("=" * 70) + log_info("=" * 70) while True: try: @@ -276,7 +196,7 @@ if __name__ == "__main__": # 如果没工作,休息标准间隔(5s) time.sleep(0.1 if has_work else POLL_INTERVAL) except KeyboardInterrupt: - logger.info("=" * 70) + log_info("=" * 70) log_stop("收到停止信号,服务正在关闭...") break except Exception as e: