""" 后台任务服务 - 处理长时间的抓取任务 """ import threading import time import uuid import logging from datetime import datetime from models.database import db from services.search_service import search_service logger = logging.getLogger('task_service') class BackgroundTaskService: """后台任务服务""" def __init__(self): self.active_threads = {} # task_id -> thread self.stop_flags = {} # task_id -> bool def create_task_id(self): """生成唯一任务ID""" return f"task_{datetime.now().strftime('%Y%m%d_%H%M%S')}_{uuid.uuid4().hex[:8]}" def start_fetch_task(self, results, auto_save=False, category=None, keywords=None): """ 启动抓取任务 参数: - results: 搜索结果列表 [{title, url, ...}, ...] - auto_save: 是否自动保存到内容库 - category: 分类(可选) - keywords: 关键词列表(可选) """ task_id = self.create_task_id() # 创建任务记录 db.create_task(task_id, 'fetch_urls', { 'total': len(results), 'auto_save': auto_save, 'category': category, 'keywords': keywords }) # 启动后台线程 self.stop_flags[task_id] = False thread = threading.Thread( target=self._run_fetch_task, args=(task_id, results, auto_save, category, keywords), daemon=True ) thread.start() self.active_threads[task_id] = thread logger.info(f"任务已启动: {task_id}, 共 {len(results)} 个URL") return task_id def _run_fetch_task(self, task_id, results, auto_save, category, keywords): """执行抓取任务(后台线程)""" try: # 更新状态为 running db.update_task_status(task_id, 'running', total=len(results)) processed = 0 success_count = 0 failed_count = 0 saved_count = 0 fetch_results = [] for i, result in enumerate(results): # 检查是否被停止 if self.stop_flags.get(task_id, False): db.update_task_status( task_id, 'stopped', progress=processed, result={'success': success_count, 'failed': failed_count, 'saved': saved_count} ) logger.info(f"任务已停止: {task_id}") return url = result.get('url', '') title = result.get('title', url[:50]) # 更新进度 db.update_task_status( task_id, 'running', progress=processed, current_item=title ) # 抓取内容 logger.info(f"[{task_id}] 抓取 {processed + 1}/{len(results)}: {title}") fetch_result = search_service.fetch_url_content(url) if fetch_result.get('success'): success_count += 1 content = fetch_result.get('content', '') # 自动保存 if auto_save and content: try: # 检查URL是否已存在 existing = db.search_articles(url) if not any(a.get('url') == url for a in existing): # 获取抓取到的网页标题 page_title = fetch_result.get('title', '') or title article_id = db.add_article( product_names=[], # 产品名称留空,后续手动关联 category=category or '', keywords=keywords or [], summary=content[:200], content=content, source=url, url=url, search_title=title # 搜索结果的标题 ) if article_id: saved_count += 1 logger.info(f"[{task_id}] 已保存: {title}") else: logger.info(f"[{task_id}] URL已存在,跳过: {url}") except Exception as e: logger.error(f"[{task_id}] 保存失败: {e}") fetch_results.append({ 'title': title, 'url': url, 'success': True, 'content_length': len(content), 'saved': auto_save }) else: failed_count += 1 error_msg = fetch_result.get('error', '未知错误') # 记录失败URL db.add_failed_url(url, title, error_msg, source='background_task') fetch_results.append({ 'title': title, 'url': url, 'success': False, 'error': error_msg }) logger.warning(f"[{task_id}] 抓取失败: {title} - {error_msg}") processed += 1 # 短暂延迟,避免请求过快 time.sleep(0.5) # 任务完成 db.update_task_status( task_id, 'completed', progress=processed, result={ 'total': len(results), 'success': success_count, 'failed': failed_count, 'saved': saved_count } ) logger.info(f"任务完成: {task_id}, 成功 {success_count}, 失败 {failed_count}, 已保存 {saved_count}") except Exception as e: logger.error(f"任务异常: {task_id} - {e}") db.update_task_status(task_id, 'failed', error_message=str(e)) finally: # 清理 if task_id in self.active_threads: del self.active_threads[task_id] if task_id in self.stop_flags: del self.stop_flags[task_id] def stop_task(self, task_id): """停止任务""" if task_id in self.stop_flags: self.stop_flags[task_id] = True logger.info(f"请求停止任务: {task_id}") return True return False def get_task_status(self, task_id): """获取任务状态""" return db.get_task(task_id) def get_active_tasks(self): """获取活动任务列表""" return db.get_active_tasks() # 全局任务服务实例 task_service = BackgroundTaskService()