""" RSS抓取模块:多线程并发抓取RSS源,解析并提取文章信息 支持失败自动重试、并发控制、超时保护 """ import time import requests import feedparser from concurrent.futures import ThreadPoolExecutor, as_completed from .config import RSS_FEEDS, FETCH_TIMEOUT, MAX_THREADS from .utils import parse_pub_time, extract_content from .logger import get_logger def _fetch_single_feed(feed_info): """抓取单个RSS源,返回该源的所有文章列表,内置2次重试机制""" logger = get_logger() name = feed_info['name'] url = feed_info['url'] max_retries = 2 retry_delay = 1 for retry in range(max_retries + 1): try: headers = { 'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36' } response = requests.get(url, headers=headers, timeout=FETCH_TIMEOUT) response.raise_for_status() feed = feedparser.parse(response.content) if feed.bozo != 0: logger.warning("%s: 解析警告", name) articles = [] for entry in feed.entries: pub_time = parse_pub_time(entry) content = extract_content(entry) article = { 'id': entry.get('id', entry.get('link', '')), 'title': entry.get('title', '无标题').strip(), 'link': entry.get('link', ''), 'content': content, 'published': pub_time, 'source': name, } articles.append(article) logger.info("%s: 抓取成功,获取 %d 篇文章", name, len(articles)) _record_source_success(name) return articles except Exception as e: if retry < max_retries: logger.warning("%s: 抓取失败,第%d次重试: %s", name, retry + 1, e) time.sleep(retry_delay) else: logger.error("%s: 抓取失败,已重试%d次,放弃: %s", name, max_retries, e) _record_source_failure(name) return [] def fetch_all_feeds(): """并行抓取所有RSS源,返回文章列表 采用标准线程池实现,自动管理并发,避免线程泄漏 """ logger = get_logger() logger.info("开始抓取 %d 个RSS源,并发数:%d", len(RSS_FEEDS), MAX_THREADS) all_articles = [] with ThreadPoolExecutor(max_workers=MAX_THREADS) as executor: future_to_feed = {executor.submit(_fetch_single_feed, feed): feed['name'] for feed in RSS_FEEDS} for future in as_completed(future_to_feed): feed_name = future_to_feed[future] try: articles = future.result() all_articles.extend(articles) except Exception as e: logger.error("%s: 抓取任务异常: %s", feed_name, e) logger.info("所有源抓取完成,共获取 %d 篇文章", len(all_articles)) return all_articles def _record_source_success(name): try: from .monitor import get_monitor get_monitor().record_source_result(name, True) except Exception: pass def _record_source_failure(name): try: from .monitor import get_monitor get_monitor().record_source_result(name, False) except Exception: pass