When the network is down, sync_file() logs individual errors for each file. Now run_once() detects when all files fail and emits a single summary warning instead of proceeding with redundant "no changes" logs. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
119 lines
3.6 KiB
Python
119 lines
3.6 KiB
Python
"""
|
|
Watch daemon - periodically syncs LAN Excel files and migrates changed tables to SQL Server.
|
|
"""
|
|
|
|
import argparse
|
|
import logging
|
|
import os
|
|
import time
|
|
|
|
from logger import setup_logger
|
|
from migrate import load_config, get_connection, create_schema, migrate_table
|
|
from sync import sync_all_files, determine_tables_to_migrate
|
|
|
|
logger = logging.getLogger('app')
|
|
|
|
|
|
def run_once(cfg):
|
|
"""Execute one sync + migrate cycle."""
|
|
excel_dir = os.path.join(os.path.dirname(os.path.abspath(__file__)), 'Excel')
|
|
|
|
# 1. Sync files from LAN
|
|
logger.info("开始同步 LAN 文件...")
|
|
sync_results = sync_all_files(cfg, excel_dir)
|
|
|
|
# 如果全部同步失败,大概率是网络中断
|
|
all_failed = sync_results and all(s == 'error' for s, _ in sync_results.values())
|
|
if all_failed:
|
|
logger.warning(
|
|
f"所有 LAN 源文件不可达 ({len(sync_results)} 个),可能网络中断,跳过本轮同步"
|
|
)
|
|
return
|
|
|
|
# 2. Determine which tables need migration
|
|
tables_to_migrate = determine_tables_to_migrate(cfg, sync_results)
|
|
|
|
if not tables_to_migrate:
|
|
logger.info("无文件变更,跳过迁移")
|
|
return
|
|
|
|
logger.info(f"需要迁移的表: {', '.join(sorted(tables_to_migrate))}")
|
|
|
|
# 3. Open DB connection and migrate each changed table
|
|
conn = None
|
|
try:
|
|
conn = get_connection(cfg)
|
|
cursor = conn.cursor()
|
|
create_schema(cursor, cfg['schema'])
|
|
|
|
for table_cfg in cfg['tables']:
|
|
table_name = table_cfg['table']
|
|
if table_name not in tables_to_migrate:
|
|
continue
|
|
|
|
try:
|
|
rows = migrate_table(cursor, table_cfg, cfg, use_truncate=True)
|
|
logger.info(f"迁移完成: {table_name} ({rows:,} 行)")
|
|
except Exception as e:
|
|
logger.error(f"迁移失败: {table_name} - {e}", exc_info=True)
|
|
# Rollback any uncommitted transaction for this table
|
|
try:
|
|
conn.rollback()
|
|
except Exception:
|
|
pass
|
|
|
|
cursor.close()
|
|
except Exception as e:
|
|
logger.error(f"数据库连接失败: {e}", exc_info=True)
|
|
finally:
|
|
if conn:
|
|
try:
|
|
conn.close()
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
def main():
|
|
parser = argparse.ArgumentParser(description='Excel to SQL Server Watch Daemon')
|
|
parser.add_argument('--config', default=None, help='配置文件路径 (默认: config.yaml)')
|
|
parser.add_argument('--once', action='store_true', help='执行一次后退出')
|
|
parser.add_argument('--interval', type=int, default=None,
|
|
help='覆盖配置文件中的 check_interval_minutes')
|
|
args = parser.parse_args()
|
|
|
|
# Load config
|
|
cfg = load_config(args.config)
|
|
|
|
# Setup logging
|
|
log_cfg = cfg.get('logging', {})
|
|
setup_logger(
|
|
'app',
|
|
log_dir=log_cfg.get('log_dir'),
|
|
level=log_cfg.get('level', 'INFO'),
|
|
max_keep_days=log_cfg.get('max_keep_days', 30),
|
|
)
|
|
|
|
interval = args.interval or cfg.get('check_interval_minutes', 30)
|
|
|
|
logger.info("=" * 50)
|
|
logger.info("Watch daemon 启动")
|
|
logger.info(f" 同步间隔: {interval} 分钟")
|
|
logger.info("=" * 50)
|
|
|
|
while True:
|
|
try:
|
|
run_once(cfg)
|
|
except Exception as e:
|
|
logger.error(f"运行周期异常: {e}", exc_info=True)
|
|
|
|
if args.once:
|
|
logger.info("--once 模式,退出")
|
|
break
|
|
|
|
logger.info(f"下次检查: {interval} 分钟后")
|
|
time.sleep(interval * 60)
|
|
|
|
|
|
if __name__ == '__main__':
|
|
main()
|