scheduler.py

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
调度模块
实现每日14:30自动选股、每日收盘后自动回测/迭代、非交易时间定时数据更新
"""

import os
import logging
import time
import threading
from datetime import datetime, timedelta
from apscheduler.schedulers.background import BackgroundScheduler
from dotenv import load_dotenv

# 加载环境变量
load_dotenv()

# 配置日志
logger = logging.getLogger(__name__)

class Scheduler:
    """调度类"""
    
    def __init__(self):
        """初始化调度器"""
        self.scheduler = BackgroundScheduler()
        self.is_running = False
        
        # 初始化各个模块
        from src.stock_selector import StockSelector
        self.selector = StockSelector()
        
        from src.daily_adjuster import DailyAdjuster
        self.adjuster = DailyAdjuster()
        
        from src.data_manager import DataManager
        self.data_manager = DataManager()
        
        from src.message_handler import MessageHandler
        self.message_handler = MessageHandler()
        
        # 配置调度任务
        self._configure_scheduled_jobs()
    
    def _configure_scheduled_jobs(self):
        """配置调度任务"""
        logger.info("配置调度任务...")
        
        # 每日14:30自动选股
        self.scheduler.add_job(
            self._run_stock_selection,
            'cron',
            hour=14,
            minute=30,
            day_of_week='mon-fri',
            name='每日自动选股',
            id='stock_selection'
        )
        
        # 每日15:30自动回测和参数优化
        self.scheduler.add_job(
            self._run_daily_adjustment,
            'cron',
            hour=15,
            minute=30,
            day_of_week='mon-fri',
            name='每日自动回测和优化',
            id='daily_adjustment'
        )
        
        # 每周日20:00更新历史数据
        self.scheduler.add_job(
            self._update_history_data,
            'cron',
            hour=20,
            minute=0,
            day_of_week='sun',
            name='每周日更新历史数据',
            id='update_history'
        )
        
        # 每小时检查一次程序状态
        self.scheduler.add_job(
            self._check_program_status,
            'interval',
            hours=1,
            name='程序状态检查',
            id='status_check'
        )
        
        # 每日9:30-15:00实时监控昨天选出的股票
        self.scheduler.add_job(
            self._run_real_time_monitoring,
            'cron',
            hour='9-15',
            minute='*/30',
            day_of_week='mon-fri',
            name='实时监控昨天选出的股票',
            id='real_time_monitoring'
        )
        
        logger.info("调度任务配置完成")
    
    def start(self):
        """启动调度器"""
        try:
            logger.info("启动调度器...")
            
            if self.is_running:
                logger.warning("调度器已经在运行")
                return False
            
            self.scheduler.start()
            self.is_running = True
            logger.info("调度器启动成功")
            
            # 启动监控线程
            self.monitor_thread = threading.Thread(target=self._monitor_scheduler)
            self.monitor_thread.daemon = False
            self.monitor_thread.start()
            
            return True
            
        except Exception as e:
            logger.error(f"调度器启动失败:{e}")
            self.is_running = False
            return False
    
    def stop(self):
        """停止调度器"""
        try:
            logger.info("停止调度器...")
            
            if not self.is_running:
                logger.warning("调度器已经停止")
                return False
            
            self.is_running = False
            self.scheduler.shutdown()
            logger.info("调度器停止成功")
            
            return True
            
        except Exception as e:
            logger.error(f"调度器停止失败:{e}")
            return False
    
    def _monitor_scheduler(self):
        """监控调度器状态"""
        logger.info("启动调度器监控...")
        
        while self.is_running:
            try:
                # 检查调度器是否正常运行
                if not self.scheduler.running:
                    logger.error("调度器已停止运行,正在尝试重启")
                    self.scheduler.start()
                
                # 检查是否有任务失败
                for job in self.scheduler.get_jobs():
                    try:
                        job.pause()
                        job.resume()
                    except Exception as e:
                        logger.error(f"任务 {job.name} 检查失败:{e}")
                
                time.sleep(300)  # 每5分钟检查一次
                
            except Exception as e:
                logger.error(f"调度器监控失败:{e}")
                time.sleep(60)
            finally:
                # 确保线程安全地停止
                if not self.is_running:
                    break
    
    def _run_stock_selection(self):
        """执行股票选择任务"""
        try:
            logger.info("开始执行每日自动选股任务...")
            
            # 检查是否为交易日
            today = datetime.now().strftime('%Y%m%d')
            if not self._is_trading_day(today):
                logger.warning(f"{today} 不是交易日,选股任务取消")
                return
            
            # 执行选股
            result = self.selector.run()
            if not result.empty:
                self.selector.save_results(result)
                self.selector.print_results(result)
            
            logger.info("每日自动选股任务完成")
            
        except Exception as e:
            logger.error(f"每日自动选股任务失败:{e}")
    
    def _run_daily_adjustment(self):
        """执行每日自调整任务"""
        try:
            logger.info("开始执行每日自动回测和优化任务...")
            
            # 检查是否为交易日
            today = datetime.now().strftime('%Y%m%d')
            if not self._is_trading_day(today):
                logger.warning(f"{today} 不是交易日,每日自调整任务取消")
                return
            
            # 执行每日自调整
            success = self.adjuster.run()
            
            if success:
                logger.info("每日自动回测和优化任务完成")
            else:
                logger.warning("每日自动回测和优化任务执行失败")
                
        except Exception as e:
            logger.error(f"每日自动回测和优化任务失败:{e}")
    
    def _update_history_data(self):
        """执行历史数据更新任务"""
        try:
            logger.info("开始执行历史数据更新任务...")
            
            # 计算需要更新的日期范围(最近7天)
            end_date = datetime.now().strftime('%Y%m%d')
            start_date = (datetime.now() - timedelta(days=7)).strftime('%Y%m%d')
            
            # 更新历史数据
            self.data_manager.update_history_data(start_date, end_date)
            
            logger.info("历史数据更新任务完成")
            
        except Exception as e:
            logger.error(f"历史数据更新任务失败:{e}")
    
    def _run_real_time_monitoring(self):
        """执行实时监控任务"""
        try:
            logger.info("开始执行实时监控任务...")
            
            # 检查是否为交易日
            today = datetime.now().strftime('%Y%m%d')
            if not self._is_trading_day(today):
                logger.warning(f"{today} 不是交易日,实时监控任务取消")
                return
            
            # 获取昨天的选股结果
            yesterday = (datetime.now() - timedelta(days=1)).strftime('%Y%m%d')
            yesterday_results = self.selector.get_results(yesterday)
            
            if not yesterday_results:
                logger.warning(f"{yesterday} 无选股结果,实时监控任务取消")
                return
            
            logger.info(f"发现 {yesterday} 选股结果,共 {len(yesterday_results)} 只股票")
            
            # 获取实时数据
            realtime_data = self.data_manager.get_real_time_data()
            
            if not realtime_data:
                logger.warning("未获取到实时数据,实时监控任务取消")
                return
            
            # 检查最佳卖出点
            for ts_code, name in yesterday_results.items():
                if ts_code in realtime_data:
                    current_price = realtime_data[ts_code]['price']
                    optimal_price = self._calculate_optimal_sell_price(ts_code)
                    
                    if current_price >= optimal_price:
                        logger.info(f"股票 {name} ({ts_code}) 已达到最佳卖出点")
                        self.message_handler.send_monitoring_result(
                            ts_code,
                            name,
                            current_price,
                            optimal_price,
                            "卖出"
                        )
            
            logger.info("实时监控任务完成")
            
        except Exception as e:
            logger.error(f"实时监控任务失败:{e}")
    
    def _calculate_optimal_sell_price(self, ts_code):
        """计算最佳卖出价格"""
        try:
            # 简单实现:获取股票的历史数据,计算最佳卖出价格
            # 实际应用中需要根据股票的历史数据和当前市场情况计算最佳卖出价格
            return 10.0
        
        except Exception as e:
            logger.error(f"计算最佳卖出价格失败:{e}")
            return 10.0
    
    def _check_program_status(self):
        """检查程序状态"""
        try:
            logger.info("检查程序状态...")
            
            # 检查各个模块是否正常
            modules_status = {
                '数据管理': self._check_data_manager(),
                '选股模块': self._check_stock_selector(),
                '回测模块': self._check_backtester(),
                '参数优化': self._check_optimizer(),
                '每日自调整': self._check_adjuster(),
                '消息处理': self._check_message_handler()
            }
            
            # 记录状态
            logger.info("程序状态检查结果:")
            for module, status in modules_status.items():
                logger.info(f"{module}: {'正常' if status else '异常'}")
                
            # 如果有模块异常,尝试重启
            if not all(modules_status.values()):
                logger.error("发现模块异常,尝试重启程序")
                self._restart_program()
                
        except Exception as e:
            logger.error(f"程序状态检查失败:{e}")
    
    def _check_data_manager(self):
        """检查数据管理模块"""
        try:
            return self.data_manager.test_connection()
        except Exception as e:
            logger.error(f"数据管理模块检查失败:{e}")
            return False
    
    def _check_stock_selector(self):
        """检查选股模块"""
        try:
            self.selector.test_selection()
            return True
        except Exception as e:
            logger.error(f"选股模块检查失败:{e}")
            return False
    
    def _check_backtester(self):
        """检查回测模块"""
        try:
            from src.backtester import Backtester
            backtester = Backtester()
            backtester.test_backtest()
            return True
        except Exception as e:
            logger.error(f"回测模块检查失败:{e}")
            return False
    
    def _check_optimizer(self):
        """检查参数优化模块"""
        try:
            from src.parameter_optimizer import ParameterOptimizer
            optimizer = ParameterOptimizer()
            optimizer.test_optimization()
            return True
        except Exception as e:
            logger.error(f"参数优化模块检查失败:{e}")
            return False
    
    def _check_adjuster(self):
        """检查每日自调整模块"""
        try:
            self.adjuster.test_adjustment()
            return True
        except Exception as e:
            logger.error(f"每日自调整模块检查失败:{e}")
            return False
    
    def _check_message_handler(self):
        """检查消息处理模块"""
        try:
            self.message_handler.test_message_sending()
            return True
        except Exception as e:
            logger.error(f"消息处理模块检查失败:{e}")
            return False
    
    def _is_trading_day(self, date_str):
        """检查是否为交易日"""
        try:
            from chinese_calendar import is_workday
            date = datetime.strptime(date_str, '%Y%m%d').date()
            return is_workday(date)
        except Exception as e:
            logger.error(f"交易日检查失败:{e}")
            return False
    
    def _restart_program(self):
        """重启程序(简单实现)"""
        try:
            logger.info("尝试重启程序...")
            
            # 停止调度器
            self.stop()
            
            # 重启程序
            import sys
            import subprocess
            python = sys.executable
            subprocess.call([python] + sys.argv)
            
        except Exception as e:
            logger.error(f"程序重启失败:{e}")
    
    def run_test_job(self, job_id):
        """手动运行指定任务"""
        try:
            logger.info(f"手动运行任务:{job_id}")
            
            if job_id == 'stock_selection':
                self._run_stock_selection()
            elif job_id == 'daily_adjustment':
                self._run_daily_adjustment()
            elif job_id == 'update_history':
                self._update_history_data()
            elif job_id == 'status_check':
                self._check_program_status()
            else:
                logger.warning(f"未知任务ID:{job_id}")
                
            return True
            
        except Exception as e:
            logger.error(f"手动运行任务失败:{e}")
            return False
    
    def get_scheduled_jobs(self):
        """获取所有调度任务信息"""
        try:
            jobs_info = []
            
            for job in self.scheduler.get_jobs():
                next_run_time = 'N/A'
                if hasattr(job, 'next_run_time') and job.next_run_time:
                    next_run_time = job.next_run_time.strftime('%Y-%m-%d %H:%M:%S')
                
                jobs_info.append({
                    'id': job.id,
                    'name': job.name,
                    'next_run_time': next_run_time,
                    'trigger': str(job.trigger)
                })
                
            return jobs_info
            
        except Exception as e:
            logger.error(f"获取调度任务信息失败:{e}")
            return []

if __name__ == "__main__":
    # 配置日志
    logging.basicConfig(
        level=logging.INFO,
        format="%(asctime)s - %(name)s - %(levelname)s - %(message)s"
    )
    
    # 测试调度器
    logger.info("测试调度器...")
    
    scheduler = Scheduler()
    
    # 打印调度任务信息
    logger.info("调度任务列表:")
    jobs = scheduler.get_scheduled_jobs()
    for job in jobs:
        logger.info(f"{job['name']} (ID: {job['id']}) - 下次运行:{job['next_run_time']}")
    
    # 启动调度器
    logger.info("启动调度器...")
    if scheduler.start():
        logger.info("调度器运行中...")
        
        # 保持程序运行
        try:
            while True:
                time.sleep(1)
        except KeyboardInterrupt:
            logger.info("调度器停止")
            scheduler.stop()
    else:
        logger.error("调度器启动失败")