"""纯 worker 模式(无 web)— 启动调度器 + 健康监控。 用法: python -m app.worker # 独立 worker 进程 python -c "from app.entrypoints.worker import start_scheduler_thread; start_scheduler_thread()" # 在 web 进程内 inline 跑 """ from __future__ import annotations import json import sys import threading import time from pathlib import Path _PROJECT_ROOT = Path(__file__).resolve().parent.parent if str(_PROJECT_ROOT) not in sys.path: sys.path.insert(0, str(_PROJECT_ROOT)) from app.core.config import settings from app.core.db import ops as db_ops from app.core.utils.logging import get_logger, setup_logging setup_logging() logger = get_logger("worker") def start_scheduler_thread() -> None: """在当前进程内启动 scheduler + 健康监控(daemon thread)。 给 service_run.sh 用:在 uvicorn 同进程跑后台 scheduler,避免 docker 一个容器跑两个进程(supervisord / s6 那种方案)。 幂等:多次调用只启动一次。 2026-07-12: 增加 RUNTIME_MODE 感知。systemd 模式下不启动 scheduler, 因为同步由 systemd timer 触发,避免双调度器重叠。 """ logger.info(f"[worker] runtime_mode={settings.runtime_mode}") if settings.runtime_mode == "systemd": logger.info("[worker] systemd 模式:不启动进程内 scheduler,同步由 systemd timer 触发") return # 1. 注册数据源 + seed config from app.core.datasource.registry import ( build_default_registry, seed_datasource_configs, start_health_monitor, ) build_default_registry() seed_datasource_configs() # 2. seed sync registry + 恢复卡死任务 from app.core.sync.registry import recover_interrupted_syncs, seed_sync_registry seed_sync_registry() n = recover_interrupted_syncs() if n > 0: logger.warning(f"[worker] 恢复了 {n} 个中断的同步任务") # 3. seed 节假日 if settings.holidays_list: db_ops.upsert_config( "trading_calendar_holidays", json.dumps(settings.holidays_list, ensure_ascii=False), category="general", description="A 股休市日", ) # 4. seed 默认计划任务 if settings.scheduler_auto_seed: from app.core.scheduler.scheduler import seed_schedule_configs seed_schedule_configs() # 5. 注册 sync job + 启动调度器 + 健康监控 from app.core.scheduler.scheduler import ( register_sync_jobs, start_scheduler, ) register_sync_jobs() start_scheduler() start_health_monitor() logger.info("[worker] scheduler + health monitor 已 inline 启动") def main(): logger.info("[worker] market_data_sync worker 启动") logger.info(f"[worker] runtime_mode={settings.runtime_mode}") if settings.runtime_mode == "systemd": logger.info("[worker] systemd 模式:不启动进程内 scheduler,同步由 systemd timer 触发") logger.info("[worker] 仅启动健康监控(数据源心跳),按 Ctrl+C 退出") from app.core.datasource.registry import build_default_registry, seed_datasource_configs, start_health_monitor build_default_registry() seed_datasource_configs() start_health_monitor() try: while True: time.sleep(60) except KeyboardInterrupt: logger.info("[worker] 收到 SIGINT,退出") sys.exit(0) return # 1. 注册数据源 + seed config from app.core.datasource.registry import build_default_registry, seed_datasource_configs build_default_registry() seed_datasource_configs() # 2. seed sync registry + 恢复卡死任务 from app.core.sync.registry import recover_interrupted_syncs, seed_sync_registry seed_sync_registry() n = recover_interrupted_syncs() if n > 0: logger.warning(f"[worker] 恢复了 {n} 个中断的同步任务") # 3. seed 节假日 if settings.holidays_list: db_ops.upsert_config( "trading_calendar_holidays", json.dumps(settings.holidays_list, ensure_ascii=False), category="general", description="A 股休市日", ) # 4. seed 默认计划任务 if settings.scheduler_auto_seed: from app.core.scheduler.scheduler import seed_schedule_configs seed_schedule_configs() # 5. 注册 sync job + 启动调度器 + 健康监控 from app.core.scheduler.scheduler import register_sync_jobs, start_scheduler from app.core.datasource.registry import start_health_monitor register_sync_jobs() start_scheduler() start_health_monitor() logger.info("[worker] 已启动调度器 + 健康监控,按 Ctrl+C 退出") try: while True: time.sleep(60) except KeyboardInterrupt: logger.info("[worker] 收到 SIGINT,退出") from app.core.scheduler.scheduler import stop_scheduler stop_scheduler() sys.exit(0) if __name__ == "__main__": main()