init
This commit is contained in:
commit
e0af97ac7f
65 files changed
+7366
No files matched your search
@@ -0,0 +1,43 @@
|
||||
from apscheduler.schedulers.background import BackgroundScheduler
|
||||
from apscheduler.executors.pool import ThreadPoolExecutor, ProcessPoolExecutor
|
||||
from apscheduler.jobstores.memory import MemoryJobStore
|
||||
import pytz
|
||||
|
||||
from app.services.sites.factory import CrawlerRegister
|
||||
from app.utils.logger import log
|
||||
from app.core.config import get_scheduler_config
|
||||
|
||||
# 创建爬虫工厂
|
||||
crawler_factory = CrawlerRegister().register()
|
||||
|
||||
# 获取调度器配置
|
||||
scheduler_config = get_scheduler_config()
|
||||
|
||||
# 配置调度器
|
||||
jobstores = {
|
||||
'default': MemoryJobStore()
|
||||
}
|
||||
|
||||
executors = {
|
||||
'default': ThreadPoolExecutor(scheduler_config.thread_pool_size),
|
||||
'processpool': ProcessPoolExecutor(scheduler_config.process_pool_size)
|
||||
}
|
||||
|
||||
job_defaults = {
|
||||
'coalesce': scheduler_config.coalesce,
|
||||
'max_instances': scheduler_config.max_instances,
|
||||
'misfire_grace_time': scheduler_config.misfire_grace_time,
|
||||
}
|
||||
|
||||
# 创建并配置调度器
|
||||
_scheduler = BackgroundScheduler(
|
||||
jobstores=jobstores,
|
||||
executors=executors,
|
||||
job_defaults=job_defaults,
|
||||
timezone=pytz.timezone(scheduler_config.timezone)
|
||||
)
|
||||
|
||||
# 启动调度器
|
||||
_scheduler.start()
|
||||
|
||||
log.info(f"Scheduler started with timezone: {scheduler_config.timezone}")
|
||||
@@ -0,0 +1,121 @@
|
||||
import threading
|
||||
import time
|
||||
import os
|
||||
from selenium import webdriver
|
||||
from selenium.webdriver.chrome.service import Service
|
||||
from webdriver_manager.chrome import ChromeDriverManager
|
||||
from app.utils.logger import log
|
||||
|
||||
class BrowserManager:
|
||||
"""浏览器管理器,提供共享的Chrome浏览器实例"""
|
||||
_instance = None
|
||||
_lock = threading.Lock()
|
||||
_driver = None
|
||||
_driver_path = None
|
||||
_last_activity = 0
|
||||
_max_idle_time = 1800 # 最大空闲时间(秒),默认30分钟
|
||||
|
||||
def __new__(cls, *args, **kwargs):
|
||||
"""单例模式实现"""
|
||||
if cls._instance is None:
|
||||
with cls._lock:
|
||||
if cls._instance is None:
|
||||
cls._instance = super(BrowserManager, cls).__new__(cls)
|
||||
cls._instance._init_driver_path()
|
||||
cls._instance._start_idle_monitor()
|
||||
return cls._instance
|
||||
|
||||
def _init_driver_path(self):
|
||||
"""初始化ChromeDriver路径"""
|
||||
try:
|
||||
self._driver_path = ChromeDriverManager().install()
|
||||
log.info(f"ChromeDriver已安装: {self._driver_path}")
|
||||
except Exception as e:
|
||||
log.error(f"ChromeDriver安装失败: {str(e)}")
|
||||
raise
|
||||
|
||||
def _start_idle_monitor(self):
|
||||
"""启动空闲监控线程"""
|
||||
def monitor():
|
||||
while True:
|
||||
time.sleep(60) # 每分钟检查一次
|
||||
try:
|
||||
with self._lock:
|
||||
if self._driver is not None:
|
||||
current_time = time.time()
|
||||
if current_time - self._last_activity > self._max_idle_time:
|
||||
log.info(f"浏览器空闲超过{self._max_idle_time}秒,释放资源")
|
||||
self._quit_driver()
|
||||
except Exception as e:
|
||||
log.error(f"浏览器监控线程异常: {str(e)}")
|
||||
|
||||
monitor_thread = threading.Thread(target=monitor, daemon=True)
|
||||
monitor_thread.start()
|
||||
log.info("浏览器空闲监控线程已启动")
|
||||
|
||||
def get_driver(self):
|
||||
"""获取Chrome浏览器实例"""
|
||||
with self._lock:
|
||||
self._last_activity = time.time()
|
||||
if self._driver is None:
|
||||
self._create_driver()
|
||||
return self._driver
|
||||
|
||||
def _create_driver(self):
|
||||
"""创建新的Chrome浏览器实例"""
|
||||
log.info("创建新的Chrome浏览器实例")
|
||||
options = webdriver.ChromeOptions()
|
||||
# 基本配置(无头模式)
|
||||
options.add_argument("--headless")
|
||||
options.add_argument("--disable-gpu")
|
||||
options.add_argument("--no-sandbox")
|
||||
# 内存优化配置
|
||||
options.add_argument("--disable-dev-shm-usage")
|
||||
options.add_argument("--disable-extensions")
|
||||
options.add_argument("--disable-application-cache")
|
||||
options.add_argument("--js-flags=--expose-gc")
|
||||
options.add_argument("--memory-pressure-off")
|
||||
options.add_argument("--disable-default-apps")
|
||||
# 日志级别
|
||||
options.add_argument("--log-level=3")
|
||||
|
||||
self._driver = webdriver.Chrome(
|
||||
service=Service(self._driver_path),
|
||||
options=options
|
||||
)
|
||||
self._driver.set_page_load_timeout(30)
|
||||
|
||||
def _quit_driver(self):
|
||||
"""关闭浏览器实例"""
|
||||
if self._driver:
|
||||
try:
|
||||
self._driver.quit()
|
||||
log.info("浏览器实例已关闭")
|
||||
except Exception as e:
|
||||
log.error(f"关闭浏览器实例出错: {str(e)}")
|
||||
finally:
|
||||
self._driver = None
|
||||
|
||||
def release_driver(self):
|
||||
"""使用完毕后标记为活动状态"""
|
||||
with self._lock:
|
||||
self._last_activity = time.time()
|
||||
|
||||
def get_page_content(self, url, wait_time=5):
|
||||
"""获取指定URL的页面内容,并自动处理浏览器"""
|
||||
driver = self.get_driver()
|
||||
try:
|
||||
driver.get(url)
|
||||
time.sleep(wait_time) # 等待页面加载
|
||||
page_source = driver.page_source
|
||||
self.release_driver()
|
||||
return page_source, driver
|
||||
except Exception as e:
|
||||
log.error(f"获取页面内容失败: {str(e)}")
|
||||
self.release_driver()
|
||||
raise
|
||||
|
||||
def shutdown(self):
|
||||
"""关闭浏览器管理器"""
|
||||
with self._lock:
|
||||
self._quit_driver()
|
||||
@@ -0,0 +1,240 @@
|
||||
import time
|
||||
import traceback
|
||||
import threading
|
||||
from datetime import datetime
|
||||
from functools import wraps
|
||||
import pytz
|
||||
import signal
|
||||
from typing import List, Dict, Any, Optional, Callable
|
||||
|
||||
from app.services import crawler_factory, _scheduler
|
||||
from app.utils.logger import log
|
||||
from app.core import db, cache
|
||||
from app.core.config import get_crawler_config
|
||||
from app.utils.notification import notification_manager
|
||||
|
||||
# 获取爬虫配置
|
||||
crawler_config = get_crawler_config()
|
||||
|
||||
# 配置常量
|
||||
CRAWLER_INTERVAL = crawler_config.interval
|
||||
CRAWLER_TIMEOUT = crawler_config.timeout
|
||||
MAX_RETRY_COUNT = crawler_config.max_retry_count
|
||||
SHANGHAI_TZ = pytz.timezone('Asia/Shanghai')
|
||||
|
||||
class CrawlerTimeoutError(Exception):
|
||||
"""爬虫超时异常"""
|
||||
pass
|
||||
|
||||
def timeout_handler(func: Callable, timeout: int = CRAWLER_TIMEOUT) -> Callable:
|
||||
"""超时处理装饰器,支持Unix信号和线程两种实现"""
|
||||
@wraps(func)
|
||||
def wrapper(*args, **kwargs):
|
||||
# 线程实现的超时机制
|
||||
result = [None]
|
||||
exception = [None]
|
||||
completed = [False]
|
||||
|
||||
def target():
|
||||
try:
|
||||
result[0] = func(*args, **kwargs)
|
||||
except Exception as e:
|
||||
exception[0] = e
|
||||
finally:
|
||||
completed[0] = True
|
||||
|
||||
thread = threading.Thread(target=target)
|
||||
thread.daemon = True
|
||||
thread.start()
|
||||
thread.join(timeout)
|
||||
|
||||
if not completed[0]:
|
||||
error_msg = f"Function {func.__name__} timed out after {timeout} seconds"
|
||||
log.error(error_msg)
|
||||
raise CrawlerTimeoutError(error_msg)
|
||||
|
||||
if exception[0]:
|
||||
log.error(f"Function {func.__name__} raised an exception: {exception[0]}")
|
||||
raise exception[0]
|
||||
|
||||
return result[0]
|
||||
return wrapper
|
||||
|
||||
def safe_fetch(crawler_name: str, crawler, date_str: str, is_retry: bool = False) -> List[Dict[str, Any]]:
|
||||
"""安全地执行爬虫抓取,处理异常并返回结果"""
|
||||
try:
|
||||
news_list = crawler.fetch(date_str)
|
||||
if news_list and len(news_list) > 0:
|
||||
cache_key = f"crawler:{crawler_name}:{date_str}"
|
||||
cache.set_cache(key=cache_key, value=news_list, expire=0)
|
||||
|
||||
log.info(f"{crawler_name} fetch success, {len(news_list)} news fetched")
|
||||
return news_list
|
||||
else:
|
||||
log.info(f"{'Second time ' if is_retry else ''}crawler {crawler_name} failed. 0 news fetched")
|
||||
return []
|
||||
except Exception as e:
|
||||
error_msg = traceback.format_exc()
|
||||
log.error(f"{'Second time ' if is_retry else ''}crawler {crawler_name} error: {error_msg}")
|
||||
|
||||
# 发送钉钉通知
|
||||
try:
|
||||
notification_manager.notify_crawler_error(
|
||||
crawler_name=crawler_name,
|
||||
error_msg=str(e),
|
||||
date_str=date_str,
|
||||
is_retry=is_retry
|
||||
)
|
||||
except Exception as notify_error:
|
||||
log.error(f"Failed to send notification for crawler {crawler_name}: {notify_error}")
|
||||
|
||||
return []
|
||||
|
||||
def run_data_analysis(date_str: str):
|
||||
"""执行数据分析并缓存结果"""
|
||||
log.info(f"Starting data analysis for date {date_str}")
|
||||
try:
|
||||
# 导入分析模块(在这里导入避免循环依赖)
|
||||
from app.analysis.trend_analyzer import TrendAnalyzer
|
||||
from app.analysis.predictor import TrendPredictor
|
||||
|
||||
# 创建分析器实例
|
||||
analyzer = TrendAnalyzer()
|
||||
predictor = TrendPredictor()
|
||||
|
||||
# 1. 生成关键词云图数据并缓存
|
||||
log.info("Generating keyword cloud data...")
|
||||
analyzer.get_keyword_cloud(date_str, refresh=True)
|
||||
|
||||
# 2. 生成热点聚合分析数据并缓存
|
||||
log.info("Generating trend analysis data...")
|
||||
analyzer.get_analysis(date_str, analysis_type="main")
|
||||
|
||||
# 3. 生成跨平台热点分析数据并缓存
|
||||
log.info("Generating cross-platform analysis data...")
|
||||
analyzer.get_cross_platform_analysis(date_str, refresh=True)
|
||||
|
||||
# 4. 生成热点趋势预测数据并缓存
|
||||
log.info("Generating trend prediction data...")
|
||||
predictor.get_prediction(date_str)
|
||||
|
||||
# 5. 生成平台对比分析数据并缓存
|
||||
log.info("Generating platform comparison data...")
|
||||
analyzer.get_platform_comparison(date_str)
|
||||
|
||||
# 6. 生成高级分析数据并缓存
|
||||
log.info("Generating advanced analysis data...")
|
||||
analyzer.get_advanced_analysis(date_str, refresh=True)
|
||||
|
||||
# 7. 生成数据可视化分析数据并缓存
|
||||
log.info("Generating data visualization analysis...")
|
||||
analyzer.get_data_visualization(date_str, refresh=True)
|
||||
|
||||
# 8. 生成趋势预测分析数据并缓存
|
||||
log.info("Generating trend forecast data...")
|
||||
analyzer.get_trend_forecast(date_str, refresh=True)
|
||||
|
||||
log.info(f"All data analysis completed for date {date_str}")
|
||||
except Exception as e:
|
||||
error_msg = traceback.format_exc()
|
||||
log.error(f"Error during data analysis: {str(e)}")
|
||||
log.error(error_msg)
|
||||
|
||||
# 发送数据分析异常通知
|
||||
try:
|
||||
notification_manager.notify_analysis_error(
|
||||
error_msg=str(e),
|
||||
date_str=date_str
|
||||
)
|
||||
except Exception as notify_error:
|
||||
log.error(f"Failed to send analysis error notification: {notify_error}")
|
||||
|
||||
@_scheduler.scheduled_job('interval', id='crawlers_logic', seconds=CRAWLER_INTERVAL,
|
||||
max_instances=crawler_config.max_instances,
|
||||
misfire_grace_time=crawler_config.misfire_grace_time)
|
||||
def crawlers_logic():
|
||||
"""爬虫主逻辑,包含超时保护和错误处理"""
|
||||
|
||||
@timeout_handler
|
||||
def crawler_work():
|
||||
now_time = datetime.now(SHANGHAI_TZ)
|
||||
date_str = now_time.strftime("%Y-%m-%d")
|
||||
log.info(f"Starting crawler job at {now_time.strftime('%Y-%m-%d %H:%M:%S')}")
|
||||
|
||||
retry_crawler = []
|
||||
success_count = 0
|
||||
failed_crawlers = []
|
||||
|
||||
for crawler_name, crawler in crawler_factory.items():
|
||||
news_list = safe_fetch(crawler_name, crawler, date_str)
|
||||
if news_list:
|
||||
success_count += 1
|
||||
else:
|
||||
retry_crawler.append(crawler_name)
|
||||
failed_crawlers.append(crawler_name)
|
||||
|
||||
# 第二轮爬取(重试失败的爬虫)
|
||||
if retry_crawler:
|
||||
log.info(f"Retrying {len(retry_crawler)} failed crawlers")
|
||||
retry_failed = []
|
||||
for crawler_name in retry_crawler:
|
||||
news_list = safe_fetch(crawler_name, crawler_factory[crawler_name], date_str, is_retry=True)
|
||||
if news_list:
|
||||
success_count += 1
|
||||
# 从失败列表中移除成功的爬虫
|
||||
if crawler_name in failed_crawlers:
|
||||
failed_crawlers.remove(crawler_name)
|
||||
else:
|
||||
retry_failed.append(crawler_name)
|
||||
|
||||
# 记录完成时间
|
||||
end_time = datetime.now(SHANGHAI_TZ)
|
||||
duration = (end_time - now_time).total_seconds()
|
||||
log.info(f"Crawler job finished at {end_time.strftime('%Y-%m-%d %H:%M:%S')}, "
|
||||
f"duration: {duration:.2f}s, success: {success_count}/{len(crawler_factory)}")
|
||||
|
||||
# 发送通知
|
||||
try:
|
||||
notification_manager.notify_crawler_summary(
|
||||
success_count=success_count,
|
||||
total_count=len(crawler_factory),
|
||||
failed_crawlers=failed_crawlers,
|
||||
duration=duration,
|
||||
date_str=date_str
|
||||
)
|
||||
except Exception as notify_error:
|
||||
log.error(f"Failed to send crawler notification: {notify_error}")
|
||||
|
||||
# 爬取完成后执行数据分析
|
||||
log.info("Crawler job completed, starting data analysis...")
|
||||
# 使用新线程执行分析,避免阻塞主线程
|
||||
threading.Thread(target=run_data_analysis, args=(date_str,), daemon=True).start()
|
||||
|
||||
return success_count
|
||||
|
||||
try:
|
||||
return crawler_work()
|
||||
except CrawlerTimeoutError as e:
|
||||
log.error(f"Crawler job timeout: {str(e)}")
|
||||
# 发送超时通知
|
||||
try:
|
||||
notification_manager.notify_crawler_timeout(
|
||||
timeout_seconds=CRAWLER_TIMEOUT,
|
||||
date_str=date_str
|
||||
)
|
||||
except Exception as notify_error:
|
||||
log.error(f"Failed to send timeout notification: {notify_error}")
|
||||
return 0
|
||||
except Exception as e:
|
||||
log.error(f"Crawler job error: {str(e)}")
|
||||
log.error(traceback.format_exc())
|
||||
# 发送通用异常通知
|
||||
try:
|
||||
notification_manager.notify_crawler_error(
|
||||
crawler_name="crawler_job",
|
||||
error_msg=str(e),
|
||||
date_str=date_str
|
||||
)
|
||||
except Exception as notify_error:
|
||||
log.error(f"Failed to send error notification: {notify_error}")
|
||||
return 0
|
||||
Loaded 3 of 65 files, more files were not shown because too many files have changed in this diff.
Show more
Reference in new issue
Block a user